Add bootstrap relay command · Entire
Add bootstrap relay command
4f2b618→main·
Soph·3mo ago·4 files·+477 added/-1 removed
Sessions
0bdb781a9478 View transcript
Changes
4
cmd/git-sync
Mmain.go+68/-1
internal/syncer
Mintegration_test.go+68
Mprotocol_v2.go+111
Msyncer.go+230
29 unmodified lines
30
31
32
33
34
35
36
37
76 unmodified lines
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
160 unmodified lines
345
346
347
281
348
349
350
351
29 unmodified lines
return runSyncLike(ctx, "sync", args[1:], false)
case "plan":
return runSyncLike(ctx, "plan", args[1:], true)
case "bootstrap":
return runBootstrap(ctx, args[1:])
case "probe":
return runProbe(ctx, args[1:])
case "fetch":
76 unmodified lines
return nil
}
func runBootstrap(ctx context.Context, args []string) error {
fs := flag.NewFlagSet("bootstrap", flag.ContinueOnError)
fs.SetOutput(os.Stderr)
cfg := syncer.Config{}
var mappings multiStringFlag
var jsonOutput bool
fs.StringVar(&cfg.Source.URL, "source-url", "", "source repository URL")
fs.StringVar(&cfg.Target.URL, "target-url", "", "target repository URL")
fs.StringVar(&cfg.Source.Token, "source-token", envOr("GITSYNC_SOURCE_TOKEN", ""), "source token/password")
fs.StringVar(&cfg.Target.Token, "target-token", envOr("GITSYNC_TARGET_TOKEN", ""), "target token/password")
fs.StringVar(&cfg.Source.Username, "source-username", envOr("GITSYNC_SOURCE_USERNAME", "git"), "source basic auth username")
fs.StringVar(&cfg.Target.Username, "target-username", envOr("GITSYNC_TARGET_USERNAME", "git"), "target basic auth username")
fs.StringVar(&cfg.Source.BearerToken, "source-bearer-token", envOr("GITSYNC_SOURCE_BEARER_TOKEN", ""), "source bearer token")
fs.StringVar(&cfg.Target.BearerToken, "target-bearer-token", envOr("GITSYNC_TARGET_BEARER_TOKEN", ""), "target bearer token")
branches := fs.String("branch", "", "comma-separated branch list; default is all source branches")
fs.Var(&mappings, "map", "ref mapping in src:dst form; short names map branches, full refs map exact refs")
fs.BoolVar(&cfg.IncludeTags, "tags", false, "mirror tags")
fs.BoolVar(&cfg.ShowStats, "stats", false, "print transfer statistics")
fs.BoolVar(&jsonOutput, "json", false, "print JSON output")
fs.StringVar(&cfg.ProtocolMode, "protocol", envOr("GITSYNC_PROTOCOL", "auto"), "protocol mode: auto, v1, or v2")
fs.BoolVar(&cfg.Verbose, "v", false, "verbose logging")
if err := fs.Parse(args); err != nil {
return err
}
positional := fs.Args()
if cfg.Source.URL == "" && len(positional) > 0 {
cfg.Source.URL = positional[0]
}
if cfg.Target.URL == "" && len(positional) > 1 {
cfg.Target.URL = positional[1]
}
if len(positional) > 2 {
return usageError("too many positional arguments")
}
if *branches != "" {
cfg.Branches = splitCSV(*branches)
}
for _, raw := range mappings {
mapping, err := parseMapping(raw)
if err != nil {
return err
}
cfg.Mappings = append(cfg.Mappings, mapping)
}
if cfg.Source.URL == "" || cfg.Target.URL == "" {
return usageError("bootstrap requires source and target repository URLs")
}
result, err := syncer.Bootstrap(ctx, cfg)
if err != nil {
return err
}
printOutput(jsonOutput, result)
return nil
}
func runProbe(ctx context.Context, args []string) error {
fs := flag.NewFlagSet("probe", flag.ContinueOnError)
fs.SetOutput(os.Stderr)
160 unmodified lines
}
func usageError(message string) error {
usage := "usage:\n git-sync sync [flags] <source-url> <target-url>\n git-sync plan [flags] <source-url> <target-url>\n git-sync bootstrap [flags] <source-url> <target-url>\n git-sync probe [flags] <source-url> [target-url]\n git-sync fetch [flags] <source-url>\n\nsync/plan flags:\n --branch main,dev\n --map main:stable\n --tags\n --force\n --prune\n --stats\n --json\n --protocol auto|v1|v2\n --source-token ...\n --target-token ...\n --source-username git\n --target-username git\n --source-bearer-token ...\n --target-bearer-token ...\n -v\n\nbootstrap flags:\n --branch main,dev\n --map main:stable\n --tags\n --stats\n --json\n --protocol auto|v1|v2\n --source-token ...\n --target-token ...\n --source-username git\n --target-username git\n --source-bearer-token ...\n --target-bearer-token ...\n -v\n\nprobe flags:\n --tags\n --stats\n --json\n --protocol auto|v1|v2\n --source-token ...\n --source-username git\n --source-bearer-token ...\n --target-token ...\n --target-username git\n --target-bearer-token ...\n\nfetch flags:\n --branch main,dev\n --tags\n --stats\n --json\n --protocol auto|v1|v2\n --have-ref main\n --have <hash>\n --source-token ...\n --source-username git\n --source-bearer-token ...\n"
if message == "" {
return errors.New(strings.TrimSpace(usage))
}
}
Mcmd/git-sync/main.go+68/-1
72 unmodified lines
func TestBootstrap_IntegrationInitialSyncToEmptyTarget(t *testing.T) {
sourceRepo, sourceFS := newSourceRepo(t)
makeCommits(t, sourceRepo, sourceFS, 4)
targetRepo, err := git.Init(memory.NewStorage(), nil)
if err != nil {
t.Fatalf("init target repo: %v", err)
}
sourceServer := newSmartHTTPRepoServerV2(t, sourceRepo)
targetServer := newSmartHTTPRepoServer(t, targetRepo)
defer sourceServer.Close()
defer targetServer.Close()
result, err := Bootstrap(context.Background(), Config{
Source: Endpoint{URL: sourceServer.RepoURL()},
Target: Endpoint{URL: targetServer.RepoURL()},
ProtocolMode: protocolModeAuto,
ShowStats: true,
})
if err != nil {
t.Fatalf("bootstrap failed: %v", err)
}
if result.Pushed != 1 || result.Blocked != 0 || len(result.Plans) != 1 {
t.Fatalf("unexpected result: %+v", result)
}
if result.Plans[0].Action != ActionCreate {
t.Fatalf("expected create plan, got %+v", result.Plans[0])
}
assertHeadsMatch(t, sourceRepo, targetRepo, testBranch)
if sourceServer.BytesOut(serviceUploadPack, metricPack) == 0 {
t.Fatalf("expected source upload-pack response bytes")
}
if targetServer.Count(serviceReceivePack, metricPack) != 1 {
t.Fatalf("expected one receive-pack POST, got %d", targetServer.Count(serviceReceivePack, metricPack))
}
}
func TestBootstrap_IntegrationFailsWhenTargetRefExists(t *testing.T) {
sourceRepo, sourceFS := newSourceRepo(t)
makeCommits(t, sourceRepo, sourceFS, 2)
targetRepo, targetFS := newSourceRepo(t)
makeCommits(t, targetRepo, targetFS, 1)
_, err := Bootstrap(context.Background(), Config{
Source: Endpoint{URL: sourceServer.RepoURL()},
Target: Endpoint{URL: targetServer.RepoURL()},
ProtocolMode: protocolModeAuto,
})
if err == nil {
t.Fatalf("expected bootstrap failure when target ref exists")
}
if !strings.Contains(err.Error(), "already exists") {
t.Fatalf("expected existing-ref error, got %v", err)
}
if targetServer.Count(serviceReceivePack, metricPack) != 0 {
t.Fatalf("expected no receive-pack POSTs, got %d", targetServer.Count(serviceReceivePack, metricPack))
}
}
func TestRun_IntegrationResyncFetchesLessFromSource(t *testing.T) {
sourceRepo, sourceFS := newSourceRepo(t)
makeCommits(t, sourceRepo, sourceFS, 10)
Minternal/syncer/integration_test.go+68
262 unmodified lines
func (s *sourceRefService) FetchPack(ctx context.Context, conn *transportConn, desired map[plumbing.ReferenceName]desiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash) (io.ReadCloser, error) { switch s.protocol { case protocolModeV2: return fetchSourcePackV2(ctx, conn, s.v2, desired, targetRefs) case protocolModeV1: return fetchSourcePackV1(ctx, conn, s.v1, desired, targetRefs) default: return nil, fmt.Errorf("unsupported source protocol %q", s.protocol) } }
func sourceCapabilities(s *sourceRefService) []string { switch s.protocol { case protocolModeV2: 262 unmodified lines
}
func fetchSourcePackV2( ctx context.Context, conn *transportConn, adv *v2CapabilityAdvertisement, desired map[plumbing.ReferenceName]desiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, ) (io.ReadCloser, error) { body, wants, haves, err := sourceFetchRequestV2(adv, desired, targetRefs) if err != nil { return nil, err } if wants == 0 { return nil, git.NoErrAlreadyUpToDate } conn.stats.addWantsHaves("source upload-pack", wants, haves)
reader, err := postRPCStreamWithPhase(ctx, conn, transport.UploadPackServiceName, body, true, "upload-pack fetch") if err != nil { return nil, err } packReader, err := openV2FetchPackStream(reader) if err != nil { _ = reader.Close() return nil, err } return packReader, nil }
func storeV2FetchPack(repo *git.Repository, r io.Reader) error { reader := newPacketReader(r) for { 30 unmodified lines
} }
func openV2FetchPackStream(body io.ReadCloser) (io.ReadCloser, error) { reader := newPacketReader(body) for { kind, payload, err := reader.ReadPacket() if err != nil { if errors.Is(err, io.EOF) { return nil, io.ErrUnexpectedEOF } return nil, fmt.Errorf("decode protocol v2 fetch response: %w", err) }
switch kind { case packetTypeFlush: return nil, io.ErrUnexpectedEOF case packetTypeDelim, packetTypeResponseEnd: continue case packetTypeData: line := string(payload) switch line { case "packfile\n": return &wrappedReadCloser{ Reader: sideband.NewDemuxer(sideband.Sideband64k, reader.Reader()), Closer: body, }, nil case "acknowledgments\n", "shallow-info\n": if err := skipV2Section(reader); err != nil { return nil, err } default: return nil, fmt.Errorf("unexpected protocol v2 fetch section %q", strings.TrimSpace(line)) } } } }
func skipV2Section(reader *packetReader) error { for { kind, _, err := reader.ReadPacket()
} }
func sourceFetchRequestV2( adv *v2CapabilityAdvertisement, desired map[plumbing.ReferenceName]desiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, ) ([]byte, int, int, error) { wants := make([]plumbing.Hash, 0, len(desired)) for _, ref := range desired { wants = append(wants, ref.SourceHash) } wants = sortedUniqueHashes(wants) haves := sortedUniqueHashes(mapsRefValues(targetRefs)) if len(wants) == 0 { return nil, 0, 0, nil }
commandArgs := make([]string, 0, len(wants)+len(haves)+4) commandArgs = append(commandArgs, "ofs-delta", "no-progress") for _, hash := range wants { commandArgs = append(commandArgs, "want "+hash.String()) } for _, hash := range haves { commandArgs = append(commandArgs, "have "+hash.String()) } commandArgs = append(commandArgs, "done")
body, err := encodeV2CommandRequest("fetch", v2RequestCapabilities(adv), commandArgs) if err != nil { return nil, 0, 0, err } return body, len(wants), len(haves), nil }
type wrappedReadCloser struct { io.Reader io.Closer }
func (r *wrappedReadCloser) Close() error { if r.closeFn == nil { return nil } return r.closeFn() }
func Bootstrap(ctx context.Context, cfg Config) (Result, error) { if cfg.ProtocolMode == "" { cfg.ProtocolMode = protocolModeAuto } if cfg.ProtocolMode != protocolModeAuto && cfg.ProtocolMode != protocolModeV1 && cfg.ProtocolMode != protocolModeV2 { return Result{}, fmt.Errorf("unsupported protocol mode %q", cfg.ProtocolMode) } if cfg.Force { return Result{}, fmt.Errorf("bootstrap does not support --force") } if cfg.Prune { return Result{}, fmt.Errorf("bootstrap does not support --prune") } if cfg.DryRun { return Result{}, fmt.Errorf("bootstrap does not support dry-run; use plan or sync") }
stats := newStats(cfg.ShowStats) sourceConn, err := newTransportConn(cfg.Source, "source", stats) if err != nil { return Result{}, fmt.Errorf("create source transport: %w", err) } targetConn, err := newTransportConn(cfg.Target, "target", stats) if err != nil { return Result{}, fmt.Errorf("create target transport: %w", err) }
sourceRefs, sourceService, err := listSourceRefs(ctx, sourceConn, cfg) if err != nil { return Result{}, fmt.Errorf("list source refs: %w", err) } targetAdv, err := advertisedRefsV1(ctx, targetConn, transport.ReceivePackServiceName) if err != nil { return Result{}, fmt.Errorf("list target refs: %w", err) } targetRefs, err := advertisedReferences(targetAdv) if err != nil { return Result{}, fmt.Errorf("decode target refs: %w", err) }
sourceRefMap := refHashMap(sourceRefs) targetRefMap := refHashMap(targetRefs)
desiredRefs, _, err := buildDesiredRefs(sourceRefMap, cfg) if err != nil { return Result{}, err } if len(desiredRefs) == 0 { return Result{}, fmt.Errorf("no source refs matched") }
plans, err := buildBootstrapPlans(desiredRefs, targetRefMap) if err != nil { return Result{}, err }
result := Result{ Plans: plans, Stats: stats.snapshot(), Protocol: sourceService.protocol, }
packReader, err := sourceService.FetchPack(ctx, sourceConn, desiredRefs, nil) if err != nil { if errors.Is(err, git.NoErrAlreadyUpToDate) { return result, nil } return result, fmt.Errorf("fetch source pack: %w", err) } defer packReader.Close()
if err := pushPackToTarget(ctx, targetConn, targetAdv, plans, packReader, cfg.Verbose); err != nil { return result, fmt.Errorf("push target refs: %w", err) }
result.Pushed = len(plans) result.Stats = stats.snapshot() return result, nil }
func buildBootstrapPlans( desired map[plumbing.ReferenceName]desiredRef, targetRefs map[plumbing.ReferenceName]plumbing.Hash, ) ([]BranchPlan, error) { targetNames := make([]plumbing.ReferenceName, 0, len(desired)) for _, want := range desired { targetNames = append(targetNames, want.TargetRef) } sort.Slice(targetNames, func(i, j int) bool { return targetNames[i] < targetNames[j] })
plans := make([]BranchPlan, 0, len(targetNames)) for _, targetRef := range targetNames { targetHash := targetRefs[targetRef] if !targetHash.IsZero() { return nil, fmt.Errorf("target ref %s already exists; use sync for non-bootstrap runs", targetRef) } want := desired[targetRef] plans = append(plans, BranchPlan{ Branch: want.Label, SourceRef: want.SourceRef, TargetRef: want.TargetRef, SourceHash: want.SourceHash, TargetHash: plumbing.ZeroHash, Kind: want.Kind, Action: ActionCreate, Reason: fmt.Sprintf("create %s at %s", want.TargetRef, shortHash(want.SourceHash)), }) } return plans, nil }