Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
133 changes: 133 additions & 0 deletions e2e/harness/multi_repo_scenario_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
}
126 changes: 102 additions & 24 deletions internal/external/command.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import (
"encoding/json"
"fmt"
"os"
"os/exec"
"strings"
"time"

"github.com/spf13/cobra"
Expand Down Expand Up @@ -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)
Expand All @@ -144,38 +146,57 @@ 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 {
return fmt.Errorf("parsing artifacts JSON: %w", err)
}
}

// 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)
Expand All @@ -185,21 +206,78 @@ 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)
}

log.Info("State committed and pushed")
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
Expand Down
Loading
Loading