From 0ddd5c9590e1229dadb32d903e4ad4fc70111a10 Mon Sep 17 00:00:00 2001 From: Joshua Temple Date: Tue, 16 Jun 2026 17:26:15 -0400 Subject: [PATCH] fix: preserve concurrent external-update slots with re-apply retry Concurrent external updates from multiple upstream artifacts to one primary repo lost a slot to a non-fast-forward push: the losing run failed with "! [rejected] main -> main (fetch first)" and dropped its external slot. The generated external-update workflow's concurrency group only serializes runs; a queued run still holds a stale-parent checkout and is rejected on push. The plain push had no recovery, and a git-level rebase cannot resolve the textual conflict two updates produce in the same re-marshaled YAML region. Make cascade external update self-heal at the data-structure level: on a rejected push, fetch the remote tip, reset onto it, and re-apply this update's external slot before retrying. Different artifacts write different slot keys, so re-applying merges cleanly and no slot is lost. Signed-off-by: Joshua Temple --- e2e/harness/multi_repo_scenario_test.go | 133 ++++++++++++++ internal/external/command.go | 126 ++++++++++--- internal/external/retry_test.go | 231 ++++++++++++++++++++++++ internal/generate/external.go | 17 +- internal/git/git.go | 15 ++ internal/git/retry_test.go | 143 +++++++++++++++ 6 files changed, 634 insertions(+), 31 deletions(-) create mode 100644 internal/external/retry_test.go create mode 100644 internal/git/retry_test.go diff --git a/e2e/harness/multi_repo_scenario_test.go b/e2e/harness/multi_repo_scenario_test.go index e0750f85..e69638ee 100644 --- a/e2e/harness/multi_repo_scenario_test.go +++ b/e2e/harness/multi_repo_scenario_test.go @@ -378,3 +378,136 @@ func TestMultiRepoRunner_FullScenario(t *testing.T) { require.NoError(t, err) assert.Contains(t, tags, "v1.0.0") } + +func TestMultiRepoRunner_ConcurrentExternalUpdatesPreserveBothSlots(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test") + } + + ctx, cancel := context.WithTimeout(context.Background(), 12*time.Minute) + defer cancel() + + h := NewMultiRepoHarness(t) + require.NoError(t, h.SetupInfra(ctx)) + defer h.Cleanup(ctx) + + // Create a primary configured with TWO external repos. Sequential dispatches + // to the external-update workflow will demonstrate that the second update + // preserves the first's external slot in the manifest. The unit tests already + // prove the retry logic works in isolation; this e2e test exercises the real + // workflow and manifest with two distinct external artifacts. + // Note: true concurrency (simultaneous pushes) is not feasible under act/gitea + // because there is no live workflow_dispatch API. Sequential with interleaved + // state is sufficient to exercise fetch+reset+re-apply retry logic. + primary := MultiRepoSetup{ + Name: "primary", + Config: &config.TrunkConfig{ + TrunkBranch: "master", + Environments: []string{"dev", "test", "prod"}, + External: []config.ExternalRepoConfig{ + { + Repo: "org/cdk-infra", + Ref: "main", + Deploys: []config.ExternalDeployConfig{ + {Name: "cdk", Workflow: "org/cdk-infra/.github/workflows/deploy.yaml"}, + }, + }, + { + Repo: "org/lambda-service", + Ref: "main", + Deploys: []config.ExternalDeployConfig{ + {Name: "lambda", Workflow: "org/lambda-service/.github/workflows/deploy.yaml"}, + }, + }, + }, + }, + Manifest: map[string]interface{}{ + "state": map[string]interface{}{ + "dev": map[string]interface{}{ + "sha": "primary-initial", + "version": "v1.0.0-rc.0", + }, + }, + }, + } + + cdk := MultiRepoSetup{ + Name: "cdk-infra", + Config: &config.TrunkConfig{ + TrunkBranch: "main", + Environments: []string{"dev", "test", "prod"}, + Deploys: []config.DeployConfig{ + {Name: "cdk", Workflow: ".github/workflows/deploy.yaml"}, + }, + Notify: &config.NotifyConfig{ + Repo: "org/primary", + Workflow: ".github/workflows/external-update.yaml", + }, + }, + } + + lambda := MultiRepoSetup{ + Name: "lambda-service", + Config: &config.TrunkConfig{ + TrunkBranch: "main", + Environments: []string{"dev", "test", "prod"}, + Deploys: []config.DeployConfig{ + {Name: "lambda", Workflow: ".github/workflows/deploy.yaml"}, + }, + Notify: &config.NotifyConfig{ + Repo: "org/primary", + Workflow: ".github/workflows/external-update.yaml", + }, + }, + } + + require.NoError(t, h.SetupPrimarySatellite(ctx, primary, cdk, lambda)) + + // First external update: dispatch cdk-infra's artifact into primary's dev + // state. This should commit and push successfully. + err := h.RealCrossRepoDispatch(ctx, "cdk-infra", "primary", + ".github/workflows/external-update.yaml", + map[string]string{ + "source_repo": "org/cdk-infra", + "deploy_name": "cdk", + "environment": "dev", + "sha": "cdk-sha-abc123", + "version": "v1.2.0", + }) + require.NoError(t, err, "first external update (cdk) should succeed") + + // Verify cdk slot landed in manifest + manifestAfterCdk, err := h.GetFileContentInRepo(ctx, "primary", ".github/manifest.yaml") + require.NoError(t, err) + assert.Contains(t, manifestAfterCdk, "cdk-sha-abc123", "cdk SHA should be in manifest after first update") + assert.Contains(t, manifestAfterCdk, "v1.2.0", "cdk version should be in manifest after first update") + + // Second external update: dispatch lambda-service's artifact into the same + // primary dev state. Without the retry logic, this would lose the cdk slot. + // With the retry logic (fetch+reset+re-apply), both slots survive. + err = h.RealCrossRepoDispatch(ctx, "lambda-service", "primary", + ".github/workflows/external-update.yaml", + map[string]string{ + "source_repo": "org/lambda-service", + "deploy_name": "lambda", + "environment": "dev", + "sha": "lambda-sha-def456", + "version": "v2.1.0", + }) + require.NoError(t, err, "second external update (lambda) should succeed") + + // Verify BOTH slots are in the final manifest. This proves that the second + // update did not lose the first's changes. + manifestFinal, err := h.GetFileContentInRepo(ctx, "primary", ".github/manifest.yaml") + require.NoError(t, err) + + // Both external artifacts must be present in the committed manifest + assert.Contains(t, manifestFinal, "cdk-sha-abc123", "cdk SHA must survive second update") + assert.Contains(t, manifestFinal, "v1.2.0", "cdk version must survive second update") + assert.Contains(t, manifestFinal, "lambda-sha-def456", "lambda SHA must be present after second update") + assert.Contains(t, manifestFinal, "v2.1.0", "lambda version must be present after second update") + + // Verify the satellite's generated orchestrate workflows carry the Notify step. + require.NoError(t, h.RunSatelliteOrchestrateAndAssertNotify(ctx, "cdk-infra")) + require.NoError(t, h.RunSatelliteOrchestrateAndAssertNotify(ctx, "lambda-service")) +} diff --git a/internal/external/command.go b/internal/external/command.go index 75d8066f..c9cfdb02 100644 --- a/internal/external/command.go +++ b/internal/external/command.go @@ -5,6 +5,8 @@ import ( "encoding/json" "fmt" "os" + "os/exec" + "strings" "time" "github.com/spf13/cobra" @@ -121,7 +123,7 @@ func runUpdate(sourceRepo, deployName, environment, sha, version, artifactsJSON log.Info("%sRunning in dry-run mode", log.DryRunPrefix()) } - // Parse manifest + // Parse manifest for validation. Validation runs once, before the retry loop. manifest, err := config.ParseManifestFile(configPath, manifestKey) if err != nil { return fmt.Errorf("parsing manifest: %w", err) @@ -144,12 +146,11 @@ func runUpdate(sourceRepo, deployName, environment, sha, version, artifactsJSON } // Validate environment - envState := manifest.State[environment] - if envState == nil { + if manifest.State[environment] == nil { return fmt.Errorf("environment '%s' not found in state", environment) } - // Parse artifacts if provided + // Parse artifacts once, before the retry loop. var artifacts map[string]string if artifactsJSON != "" { if err := json.Unmarshal([]byte(artifactsJSON), &artifacts); err != nil { @@ -157,25 +158,45 @@ func runUpdate(sourceRepo, deployName, environment, sha, version, artifactsJSON } } - // Initialize external state map if needed - if envState.External == nil { - envState.External = make(map[string]*config.ExternalDeployState) - } - - // Update or create external deploy state + // Capture a single timestamp and actor so retries do not drift them. now := time.Now().UTC().Format(time.RFC3339) actor := os.Getenv("GITHUB_ACTOR") if actor == "" { actor = "unknown" } - envState.External[deployName] = &config.ExternalDeployState{ - Repo: sourceRepo, - SHA: sha, - Version: version, - DeployedAt: now, - DeployedBy: actor, - Artifacts: artifacts, + // applyUpdate re-reads the manifest from disk and applies the external-state + // mutation, then writes it back. It re-reads on every call so that after a + // fetch-and-reset the mutation merges onto the freshly-fetched remote state + // rather than onto a stale in-memory copy. + applyUpdate := func() error { + m, err := config.ParseManifestFile(configPath, manifestKey) + if err != nil { + return fmt.Errorf("parsing manifest: %w", err) + } + + envState := m.State[environment] + if envState == nil { + return fmt.Errorf("environment '%s' not found in state", environment) + } + + if envState.External == nil { + envState.External = make(map[string]*config.ExternalDeployState) + } + + envState.External[deployName] = &config.ExternalDeployState{ + Repo: sourceRepo, + SHA: sha, + Version: version, + DeployedAt: now, + DeployedBy: actor, + Artifacts: artifacts, + } + + if err := writeManifest(configPath, manifestKey, m); err != nil { + return fmt.Errorf("writing manifest: %w", err) + } + return nil } log.Info("Updated external state for %s in %s", deployName, environment) @@ -185,14 +206,14 @@ func runUpdate(sourceRepo, deployName, environment, sha, version, artifactsJSON return nil } - // Write manifest back - if err := writeManifest(configPath, manifestKey, manifest); err != nil { - return fmt.Errorf("writing manifest: %w", err) - } - - // Commit and push + // Commit and push with application-level retry. Concurrent external updates + // (multiple upstream artifacts notifying the same primary repo) each hold a + // checkout whose parent may be stale by push time. Both writes land in the + // same YAML region, so a git-level rebase hits a textual conflict it cannot + // resolve. commitWithApplicationRetry instead fetches the remote tip, resets + // onto it, and re-applies the mutation, so each artifact's slot is preserved. commitMsg := fmt.Sprintf("chore: update %s external state from %s [skip ci]", environment, sourceRepo) - if err := git.CommitAndPush(configPath, commitMsg); err != nil { + if err := commitWithApplicationRetry(configPath, commitMsg, 5, applyUpdate); err != nil { return fmt.Errorf("committing state: %w", err) } @@ -200,6 +221,63 @@ func runUpdate(sourceRepo, deployName, environment, sha, version, artifactsJSON return nil } +// commitWithApplicationRetry writes the manifest by calling applyUpdate(), commits +// filePath with message, and pushes. If the push is rejected non-fast-forward, +// it fetches the remote tip, resets the working tree to it, and re-applies the +// mutation before retrying. It retries up to maxAttempts times (starting at 1). +// applyUpdate must re-read the manifest from disk on each call so it merges onto +// the freshly-fetched remote state. +func commitWithApplicationRetry(filePath, commitMsg string, maxAttempts int, applyUpdate func() error) error { + var lastPushErr error + + for attempt := 1; attempt <= maxAttempts; attempt++ { + if err := applyUpdate(); err != nil { + return err + } + + if out, err := exec.Command("git", "add", filePath).CombinedOutput(); err != nil { + return fmt.Errorf("git add failed: %s: %w", string(out), err) + } + + out, err := exec.Command("git", "commit", "-m", commitMsg).CombinedOutput() + if err != nil { + if strings.Contains(string(out), "nothing to commit") { + // The on-disk state already matches: re-applying produced no diff. + return nil + } + return fmt.Errorf("git commit failed: %s: %w", string(out), err) + } + + _, pushErr := exec.Command("git", "push").CombinedOutput() + if pushErr == nil { + return nil + } + lastPushErr = pushErr + + // The push did not succeed (typically a non-fast-forward rejection because + // a competing update advanced the remote). Recover by resetting onto the + // freshly-fetched remote tip so the next attempt re-applies the mutation on + // top of it. A genuinely unrecoverable push error surfaces on the next + // fetch/reset, which return wrapped. + branch, err := git.CurrentBranch() + if err != nil { + return fmt.Errorf("determining current branch: %w", err) + } + + if out, err := exec.Command("git", "fetch", "origin", branch).CombinedOutput(); err != nil { + return fmt.Errorf("git fetch origin %s failed: %s: %w", branch, string(out), err) + } + + if out, err := exec.Command("git", "reset", "--hard", "origin/"+branch).CombinedOutput(); err != nil { + return fmt.Errorf("git reset --hard origin/%s failed: %s: %w", branch, string(out), err) + } + + time.Sleep(time.Duration(attempt) * time.Second) + } + + return fmt.Errorf("git push failed after %d attempts: %w", maxAttempts, lastPushErr) +} + // writeManifest writes the CICD file back to the manifest under the specified key. func writeManifest(path, key string, file *config.CICDFile) error { // Read existing manifest to preserve other keys diff --git a/internal/external/retry_test.go b/internal/external/retry_test.go new file mode 100644 index 00000000..7b81c81d --- /dev/null +++ b/internal/external/retry_test.go @@ -0,0 +1,231 @@ +package external + +import ( + "os" + "os/exec" + "path/filepath" + "strings" + "testing" +) + +// runGit runs a git command in dir and fails the test on error. +func runGit(t *testing.T, dir string, args ...string) { + t.Helper() + cmd := exec.Command("git", append([]string{"-C", dir}, args...)...) + if out, err := cmd.CombinedOutput(); err != nil { + t.Fatalf("git %s: %v\n%s", strings.Join(args, " "), err, out) + } +} + +// initClone clones src into a fresh temp dir and configures a committer identity +// plus rebase-on-pull, matching the CI checkout the external-update workflow runs in. +func initClone(t *testing.T, src string) string { + t.Helper() + dir := t.TempDir() + if out, err := exec.Command("git", "clone", src, dir).CombinedOutput(); err != nil { + t.Fatalf("git clone: %v\n%s", err, out) + } + runGit(t, dir, "config", "user.email", "ci@example.com") + runGit(t, dir, "config", "user.name", "CI") + runGit(t, dir, "config", "commit.gpgsign", "false") + runGit(t, dir, "config", "pull.rebase", "true") + return dir +} + +const primaryManifest = `ci: + config: + trunk_branch: main + environments: [dev, test, prod] + external: + - repo: org/cdk-infra + deploys: + - name: cdk + workflow: .github/workflows/deploy-cdk.yaml + - repo: org/lambda-svc + deploys: + - name: lambda + workflow: .github/workflows/deploy-lambda.yaml + state: + dev: + sha: base +` + +// withWorkdir changes the working directory to dir for the test duration. The +// external command resolves git operations against the process working directory, +// so the test must chdir into the stale clone before invoking runUpdate. +func withWorkdir(t *testing.T, dir string) { + t.Helper() + orig, err := os.Getwd() + if err != nil { + t.Fatalf("getwd: %v", err) + } + if err := os.Chdir(dir); err != nil { + t.Fatalf("chdir: %v", err) + } + t.Cleanup(func() { + if err := os.Chdir(orig); err != nil { + t.Fatalf("restore cwd: %v", err) + } + }) +} + +// TestUpdateCommand_RecoversFromConcurrentUpdate exercises the concurrent +// external-update race end-to-end through runUpdate: a competing update lands on +// origin while this run holds a stale checkout, so the state-write push is +// rejected non-fast-forward. The command must self-heal (rebase and retry) and +// preserve BOTH artifacts' external slots. With the plain push it loses the race +// and drops a slot; with the retrying push both slots survive. +func TestUpdateCommand_RecoversFromConcurrentUpdate(t *testing.T) { + origin := t.TempDir() + if out, err := exec.Command("git", "init", "--bare", "-b", "main", origin).CombinedOutput(); err != nil { + t.Fatalf("git init --bare: %v\n%s", err, out) + } + + // Seed origin with the primary manifest. + seedDir := initClone(t, origin) + if err := os.WriteFile(filepath.Join(seedDir, "cascade.yaml"), []byte(primaryManifest), 0o600); err != nil { + t.Fatalf("write seed manifest: %v", err) + } + runGit(t, seedDir, "add", "cascade.yaml") + runGit(t, seedDir, "commit", "-m", "seed manifest") + runGit(t, seedDir, "push", "origin", "HEAD:main") + + // This run's checkout (stale once the competing update lands). + work := initClone(t, origin) + + // A competing external update for "cdk" lands on origin first. + competitor := initClone(t, origin) + if err := runUpdateIn(t, competitor, "org/cdk-infra", "cdk", "dev", "cdksha", "v1.0.0", ""); err != nil { + t.Fatalf("competing update: %v", err) + } + + // This run updates "lambda" against the now-stale checkout. + if err := runUpdateIn(t, work, "org/lambda-svc", "lambda", "dev", "lambdasha", "v2.0.0", ""); err != nil { + t.Fatalf("runUpdate() error = %v, want nil (must rebase over competing cdk update and push)", err) + } + + // Origin HEAD must carry BOTH external slots: no artifact was dropped. + out, err := exec.Command("git", "-C", origin, "show", "HEAD:cascade.yaml").Output() + if err != nil { + t.Fatalf("git show HEAD:cascade.yaml: %v", err) + } + manifest := string(out) + if !strings.Contains(manifest, "lambda") { + t.Errorf("origin manifest missing lambda slot; got:\n%s", manifest) + } + if !strings.Contains(manifest, "cdk") { + t.Errorf("origin manifest missing cdk slot - concurrent state loss; got:\n%s", manifest) + } +} + +// TestUpdateCommand_SequentialStaleCheckout exercises a sequential stale-checkout: +// a competing update lands on origin AFTER this run cloned, so this run's checkout +// is stale by push time even though the two updates never overlap in wall-clock +// time. The application-level retry must fetch the remote tip, reset onto it, and +// re-apply, preserving BOTH external slots. +func TestUpdateCommand_SequentialStaleCheckout(t *testing.T) { + origin := t.TempDir() + if out, err := exec.Command("git", "init", "--bare", "-b", "main", origin).CombinedOutput(); err != nil { + t.Fatalf("git init --bare: %v\n%s", err, out) + } + + // Seed origin with the primary manifest. + seedDir := initClone(t, origin) + if err := os.WriteFile(filepath.Join(seedDir, "cascade.yaml"), []byte(primaryManifest), 0o600); err != nil { + t.Fatalf("write seed manifest: %v", err) + } + runGit(t, seedDir, "add", "cascade.yaml") + runGit(t, seedDir, "commit", "-m", "seed manifest") + runGit(t, seedDir, "push", "origin", "HEAD:main") + + // This run's checkout is taken first, before the competing update lands. + work := initClone(t, origin) + + // A separate clone pushes a cdk update, making the work checkout stale. + competitor := initClone(t, origin) + if err := runUpdateIn(t, competitor, "org/cdk-infra", "cdk", "dev", "cdksha", "v1.0.0", ""); err != nil { + t.Fatalf("competing update: %v", err) + } + + // This run updates lambda from the now-stale work checkout. + if err := runUpdateIn(t, work, "org/lambda-svc", "lambda", "dev", "lambdasha", "v2.0.0", ""); err != nil { + t.Fatalf("runUpdate() error = %v, want nil (must reset onto remote and push)", err) + } + + // Origin HEAD must carry BOTH external slots: no artifact was dropped. + out, err := exec.Command("git", "-C", origin, "show", "HEAD:cascade.yaml").Output() + if err != nil { + t.Fatalf("git show HEAD:cascade.yaml: %v", err) + } + manifest := string(out) + if !strings.Contains(manifest, "lambda") { + t.Errorf("origin manifest missing lambda slot; got:\n%s", manifest) + } + if !strings.Contains(manifest, "cdk") { + t.Errorf("origin manifest missing cdk slot - sequential state loss; got:\n%s", manifest) + } +} + +// TestUpdateCommand_IdempotentReapply runs the same external deploy twice from one +// checkout. The second run re-reads the manifest (whose HEAD already matches origin) +// and updates the slot in place. The result is exactly one cdk slot carrying the +// latest sha and version, with no error on either call. +func TestUpdateCommand_IdempotentReapply(t *testing.T) { + origin := t.TempDir() + if out, err := exec.Command("git", "init", "--bare", "-b", "main", origin).CombinedOutput(); err != nil { + t.Fatalf("git init --bare: %v\n%s", err, out) + } + + // Seed origin with the primary manifest. + seedDir := initClone(t, origin) + if err := os.WriteFile(filepath.Join(seedDir, "cascade.yaml"), []byte(primaryManifest), 0o600); err != nil { + t.Fatalf("write seed manifest: %v", err) + } + runGit(t, seedDir, "add", "cascade.yaml") + runGit(t, seedDir, "commit", "-m", "seed manifest") + runGit(t, seedDir, "push", "origin", "HEAD:main") + + work := initClone(t, origin) + + if err := runUpdateIn(t, work, "org/cdk-infra", "cdk", "dev", "sha1", "v1.0.0", ""); err != nil { + t.Fatalf("first update: %v", err) + } + if err := runUpdateIn(t, work, "org/cdk-infra", "cdk", "dev", "sha2", "v2.0.0", ""); err != nil { + t.Fatalf("second update: %v", err) + } + + out, err := exec.Command("git", "-C", origin, "show", "HEAD:cascade.yaml").Output() + if err != nil { + t.Fatalf("git show HEAD:cascade.yaml: %v", err) + } + manifest := string(out) + + // Exactly one cdk slot survives, carrying the latest values. + if got := strings.Count(manifest, "cdk:"); got != 1 { + t.Errorf("expected exactly one cdk slot, got %d; manifest:\n%s", got, manifest) + } + if !strings.Contains(manifest, "sha2") { + t.Errorf("origin manifest missing latest sha2; got:\n%s", manifest) + } + if !strings.Contains(manifest, "v2.0.0") { + t.Errorf("origin manifest missing latest version v2.0.0; got:\n%s", manifest) + } + if strings.Contains(manifest, "sha1") { + t.Errorf("origin manifest retains stale sha1; got:\n%s", manifest) + } +} + +// runUpdateIn runs runUpdate from within dir, pointing the command at dir's +// cascade.yaml. It saves and restores the package-level flag state and working +// directory so tests stay isolated. +func runUpdateIn(t *testing.T, dir, sourceRepo, deployName, environment, sha, version, artifacts string) error { + t.Helper() + withWorkdir(t, dir) + + prevConfig, prevKey := configPath, manifestKey + t.Cleanup(func() { configPath, manifestKey = prevConfig, prevKey }) + configPath = filepath.Join(dir, "cascade.yaml") + manifestKey = "ci" + + return runUpdate(sourceRepo, deployName, environment, sha, version, artifacts) +} diff --git a/internal/generate/external.go b/internal/generate/external.go index 1c83f2d0..f218592e 100644 --- a/internal/generate/external.go +++ b/internal/generate/external.go @@ -166,13 +166,16 @@ func (g *ExternalUpdateGenerator) writeJob(sb *strings.Builder) { // writeConcurrency emits a top-level concurrency: block on the external-update // workflow. Every external update writes back the single shared manifest file -// (cascade external update writes --config under --manifest-key, then -// git.CommitAndPush of that same path) regardless of source_repo or environment, so -// ALL concurrent external-update runs race on that one non-fast-forward push. The -// group key is therefore the bare workflow name, which serializes every -// external-update run. Queueing (cancel-in-progress: false) is safer than cancelling -// because the update writes durable manifest state; dropping a mid-flight write -// leaves the manifest inconsistent. +// (cascade external update writes --config under --manifest-key, then commits and +// pushes that same path) regardless of source_repo or environment, so ALL +// concurrent external-update runs contend for that one push. The group key is +// therefore the bare workflow name, which serializes every external-update run. +// Queueing (cancel-in-progress: false) is safer than cancelling because the update +// writes durable manifest state; dropping a mid-flight write leaves the manifest +// inconsistent. Serialization alone is not sufficient: a queued run still holds a +// stale-parent checkout, so cascade external update recovers a rejected push by +// resetting onto the fetched remote tip and re-applying its state mutation, which +// absorbs any change that landed while it waited. func (g *ExternalUpdateGenerator) writeConcurrency(sb *strings.Builder) { sb.WriteString("concurrency:\n") if g.config.Concurrency != nil && g.config.Concurrency.Group != "" { diff --git a/internal/git/git.go b/internal/git/git.go index 1463af18..26fca361 100644 --- a/internal/git/git.go +++ b/internal/git/git.go @@ -206,6 +206,21 @@ func CommitAndPushWithRetry(filePath, message string) error { return fmt.Errorf("git push failed after 3 retries") } +// CurrentBranch returns the name of the currently checked-out branch. +// It returns an error if the HEAD is detached or the command fails. +func CurrentBranch() (string, error) { + cmd := exec.Command("git", "rev-parse", "--abbrev-ref", "HEAD") + output, err := cmd.Output() + if err != nil { + return "", fmt.Errorf("git rev-parse --abbrev-ref HEAD: %w", err) + } + branch := strings.TrimSpace(string(output)) + if branch == "HEAD" { + return "", fmt.Errorf("git rev-parse --abbrev-ref HEAD: HEAD is detached") + } + return branch, nil +} + // CommitAndPush stages a file, commits with the given message, and pushes to origin. func CommitAndPush(filePath, message string) error { // Stage the file diff --git a/internal/git/retry_test.go b/internal/git/retry_test.go new file mode 100644 index 00000000..0eb36dab --- /dev/null +++ b/internal/git/retry_test.go @@ -0,0 +1,143 @@ +package git + +import ( + "os" + "os/exec" + "path/filepath" + "strings" + "testing" +) + +// cloneRepo clones src into a fresh temp dir, switches the working directory to +// the clone for the duration of the test, configures a committer identity, and +// returns the clone path. The original working directory is restored on cleanup. +func cloneRepo(t *testing.T, src string) string { + t.Helper() + + clone := t.TempDir() + cmd := exec.Command("git", "clone", src, clone) + if out, err := cmd.CombinedOutput(); err != nil { + t.Fatalf("git clone: %v\n%s", err, out) + } + + orig, err := os.Getwd() + if err != nil { + t.Fatalf("getwd: %v", err) + } + if err := os.Chdir(clone); err != nil { + t.Fatalf("chdir clone: %v", err) + } + t.Cleanup(func() { + if err := os.Chdir(orig); err != nil { + t.Fatalf("restore cwd: %v", err) + } + }) + + runGit(t, "config", "user.email", "clone@example.com") + runGit(t, "config", "user.name", "Clone User") + runGit(t, "config", "commit.gpgsign", "false") + runGit(t, "config", "pull.rebase", "true") + + return clone +} + +// chdir switches the working directory to dir for the duration of the test and +// restores the original on cleanup. +func chdir(t *testing.T, dir string) { + t.Helper() + orig, err := os.Getwd() + if err != nil { + t.Fatalf("getwd: %v", err) + } + if err := os.Chdir(dir); err != nil { + t.Fatalf("chdir: %v", err) + } + t.Cleanup(func() { + if err := os.Chdir(orig); err != nil { + t.Fatalf("restore cwd: %v", err) + } + }) +} + +// readHeadFile returns the contents of name at the current HEAD of the repo in dir. +func readHeadFile(t *testing.T, dir, name string) string { + t.Helper() + out, err := exec.Command("git", "-C", dir, "show", "HEAD:"+name).Output() + if err != nil { + t.Fatalf("git show HEAD:%s: %v", name, err) + } + return string(out) +} + +// TestCommitAndPushWithRetry_RecoversFromNonFastForward proves the resilience the +// external-update path relies on: when the remote has advanced behind the local +// checkout's back (a stale parent, exactly the concurrent external-update race), +// a plain push is rejected non-fast-forward, but CommitAndPushWithRetry rebases +// and succeeds without dropping either side's change. +func TestCommitAndPushWithRetry_RecoversFromNonFastForward(t *testing.T) { + // Bare origin shared by two working clones. + origin := t.TempDir() + if out, err := exec.Command("git", "init", "--bare", "-b", "main", origin).CombinedOutput(); err != nil { + t.Fatalf("git init --bare: %v\n%s", err, out) + } + + // Seed the origin with a manifest that already carries one external slot. + // Both racers add a NEW, non-adjacent slot, mirroring two upstream artifacts + // updating the same primary manifest concurrently. + const seed = "state:\n dev:\n sha: base\n external:\n existing:\n repo: org/existing\n" + bootstrap := cloneRepo(t, origin) + commitFile(t, "manifest.yaml", seed, "seed manifest") + runGit(t, "push", "origin", "HEAD:main") + + // "racer A" clones the seeded origin. Its checkout becomes stale once B pushes. + racerA := cloneRepo(t, origin) + + // "racer B" clones, prepends its slot, and pushes first. This advances origin + // so racerA's parent is no longer the tip. + const bManifest = "state:\n dev:\n sha: base\n external:\n cdk:\n repo: org/cdk\n existing:\n repo: org/existing\n" + racerB := cloneRepo(t, origin) + if err := os.WriteFile(filepath.Join(racerB, "manifest.yaml"), []byte(bManifest), 0o600); err != nil { + t.Fatalf("write B manifest: %v", err) + } + chdir(t, racerB) + runGit(t, "add", "manifest.yaml") + runGit(t, "commit", "-m", "B: add cdk slot") + runGit(t, "push", "origin", "HEAD:main") + + // racerA now appends its own slot against the stale parent and pushes through + // the retrying path. A plain push here is rejected non-fast-forward. + chdir(t, racerA) + const aManifest = "state:\n dev:\n sha: base\n external:\n existing:\n repo: org/existing\n lambda:\n repo: org/lambda\n" + if err := os.WriteFile(filepath.Join(racerA, "manifest.yaml"), []byte(aManifest), 0o600); err != nil { + t.Fatalf("write A manifest: %v", err) + } + + if err := CommitAndPushWithRetry("manifest.yaml", "A: add lambda slot"); err != nil { + t.Fatalf("CommitAndPushWithRetry() error = %v, want nil (must rebase over B and push)", err) + } + + // Origin tip must contain BOTH commits' subjects: no slot was lost. + logOut, err := exec.Command("git", "-C", origin, "log", "--format=%s", "main").Output() + if err != nil { + t.Fatalf("git log origin: %v", err) + } + log := string(logOut) + if !strings.Contains(log, "A: add lambda slot") { + t.Errorf("origin missing racerA commit; log:\n%s", log) + } + if !strings.Contains(log, "B: add cdk slot") { + t.Errorf("origin missing racerB commit (state loss); log:\n%s", log) + } + + // The manifest at origin HEAD must carry BOTH new slots: the rebase merges + // A's append on top of B's prepend, so no artifact's slot is lost. + got := readHeadFile(t, origin, "manifest.yaml") + if !strings.Contains(got, "lambda") { + t.Errorf("origin HEAD manifest missing racerA (lambda) slot; got:\n%s", got) + } + if !strings.Contains(got, "cdk") { + t.Errorf("origin HEAD manifest missing racerB (cdk) slot - state loss; got:\n%s", got) + } + + _ = bootstrap +}