From da48de56354ef6b6378de65f13490e8cc6d13654 Mon Sep 17 00:00:00 2001 From: squarepots <46488165+squarepots@users.noreply.github.com> Date: Fri, 18 Sep 2026 00:48:51 +0800 Subject: [PATCH] Make route deployment retries recover safely Persist deployment intent, audit before retry mutation, and keep migration retries on the same recovery contract. Signed-off-by: squarepots <46488165+squarepots@users.noreply.github.com> --- internal/steward/deployment_recovery_test.go | 82 ++++++++++++++++++++ internal/steward/execution.go | 46 ++++++++--- version.txt | 2 +- 3 files changed, 120 insertions(+), 10 deletions(-) create mode 100644 internal/steward/deployment_recovery_test.go diff --git a/internal/steward/deployment_recovery_test.go b/internal/steward/deployment_recovery_test.go new file mode 100644 index 0000000..261bfaa --- /dev/null +++ b/internal/steward/deployment_recovery_test.go @@ -0,0 +1,82 @@ +package steward + +import ( + "context" + "errors" + "testing" +) + +func TestDeploymentAuditAllowsOnlyKnownRemoteStates(t *testing.T) { + route := &Route{State: "pending"} + if err := deploymentAuditAllowsMutation(route, AuditEvidence{Category: "service-missing"}); err != nil { + t.Fatalf("fresh pending Route rejected service-missing state: %v", err) + } + if err := deploymentAuditAllowsMutation(route, AuditEvidence{Category: "in-sync"}); err != nil { + t.Fatalf("pending Route rejected verified in-sync state: %v", err) + } + if err := deploymentAuditAllowsMutation(route, AuditEvidence{Category: "undetermined"}); err == nil { + t.Fatal("pending Route accepted undetermined remote state") + } + + route.State = "deploying" + if err := deploymentAuditAllowsMutation(route, AuditEvidence{Category: "service-missing"}); err != nil { + t.Fatalf("retry checkpoint rejected known missing service: %v", err) + } + + route.State = "deployed" + if err := deploymentAuditAllowsMutation(route, AuditEvidence{Category: "service-missing"}); err == nil { + t.Fatal("deployed Route accepted a missing remote service") + } +} + +func TestMarkRouteDeployingPersistsRetryCheckpoint(t *testing.T) { + state, route := healthFixture(t, "direct", false) + current := findRoute(state.Inventory, route.ID) + current.Enabled = false + current.State = "pending" + if err := state.Save(false); err != nil { + t.Fatal(err) + } + if err := markRouteDeploying(state, route.ID); err != nil { + t.Fatal(err) + } + reloaded, err := LoadState(state.PrivateDir) + if err != nil { + t.Fatal(err) + } + checkpoint := findRoute(reloaded.Inventory, route.ID) + if checkpoint == nil || checkpoint.State != "deploying" || checkpoint.Enabled { + t.Fatalf("deployment retry checkpoint was not persisted: %#v", checkpoint) + } +} + +func TestMigrationRetriesAfterUndeterminedDeploymentWithoutDuplicatingState(t *testing.T) { + state, source, input := migrationFixture(t, "direct", false) + deps := migrationTestDependencies(nil, "healthy") + originalDeploy := deps.Deploy + calls := 0 + deps.Deploy = func(ctx context.Context, state *State, routeID string) (map[string]any, error) { + calls++ + if calls == 1 { + return nil, &operationStageError{ + Stage: "remote-deployment", + StateChanged: "remote-state-undetermined", + Retry: "deploy-route", + Err: errors.New("synthetic interrupted deployment"), + } + } + return originalDeploy(ctx, state, routeID) + } + + blocked, err := migrateRouteWith(context.Background(), state, source.ID, input, deps) + if err != nil || blocked.Status != "blocked" || blocked.Phase != "replacement-prepared" || blocked.LastFailure != "replacement-deployment-failed" { + t.Fatalf("partial deployment was not checkpointed for retry: %#v err=%v", blocked, err) + } + completed, err := migrateRouteWith(context.Background(), state, source.ID, input, deps) + if err != nil || completed.Status != "complete" || calls != 2 { + t.Fatalf("migration did not resume through the deployment contract: %#v calls=%d err=%v", completed, calls, err) + } + if len(state.Inventory.Routes) != 2 || len(state.Inventory.Servers) != 2 { + t.Fatalf("migration retry duplicated desired state: routes=%d servers=%d", len(state.Inventory.Routes), len(state.Inventory.Servers)) + } +} diff --git a/internal/steward/execution.go b/internal/steward/execution.go index 8cc1874..91a7460 100644 --- a/internal/steward/execution.go +++ b/internal/steward/execution.go @@ -58,34 +58,62 @@ func deployRouteWithoutRender(ctx context.Context, state *State, routeID string) return deployRoute(ctx, state, routeID, true, false) } +func deploymentAuditAllowsMutation(route *Route, current AuditEvidence) error { + if current.Category == "in-sync" { + return nil + } + if current.Category == "service-missing" && route.State != "deployed" { + return nil + } + return errors.New("remote Route state is not safe to overwrite") +} + +func markRouteDeploying(state *State, routeID string) error { + candidate := cloneInventory(state.Inventory) + route := findRoute(candidate, routeID) + if route == nil { + return fmt.Errorf("unknown Route %q", routeID) + } + route.State = "deploying" + return saveCandidate(state, candidate, false) +} + func deployRoute(ctx context.Context, state *State, routeID string, skipClientValidation, renderClients bool) (map[string]any, error) { route := findRoute(state.Inventory, routeID) if route == nil { return nil, fmt.Errorf("unknown Route %q", routeID) } - if route.State == "deployed" { - current := AuditRoute(ctx, state, routeID) - if current.Category != "in-sync" { - return nil, errors.New("existing deployed Route has drift; refusing to overwrite unknown remote state") - } + + current := AuditRoute(ctx, state, routeID) + if err := deploymentAuditAllowsMutation(route, current); err != nil { + return nil, &operationStageError{Stage: "remote-preflight-audit", StateChanged: "remote-state-unchanged", Retry: "audit", Err: err} } + if err := markRouteDeploying(state, routeID); err != nil { + return nil, &operationStageError{Stage: "deployment-intent", StateChanged: "remote-state-unchanged", Retry: "deploy-route", Err: err} + } + route = findRoute(state.Inventory, routeID) + lines, err := performRouteOperation(ctx, state, *route, false) if err != nil { - return nil, errors.New("deterministic route operation failed") + return nil, &operationStageError{Stage: "remote-deployment", StateChanged: "remote-state-undetermined", Retry: "deploy-route", Err: errors.New("route deployment did not produce verified remote state")} } evidence := parseAuditEvidence(state.Inventory, *route, lines) if evidence.Status != "healthy" { - return nil, errors.New("route deployment completed but post-deploy audit is not healthy") + return nil, &operationStageError{Stage: "remote-deployment", StateChanged: "remote-state-undetermined", Retry: "deploy-route", Err: errors.New("route deployment completed but post-deploy audit is not healthy")} } + return adoptVerifiedRoute(state, route, evidence, skipClientValidation, renderClients) +} + +func adoptVerifiedRoute(state *State, route *Route, evidence AuditEvidence, skipClientValidation, renderClients bool) (map[string]any, error) { candidate := cloneInventory(state.Inventory) candidateRoute := findRoute(candidate, route.ID) candidateRoute.Enabled = true candidateRoute.State = "deployed" if err := saveCandidate(state, candidate, false); err != nil { - return nil, err + return nil, &operationStageError{Stage: "local-route-state", StateChanged: "route-deployed-verified", Retry: "deploy-route", Err: err} } if _, err := SetObservedRoute(state, route.ID, evidence.Status, evidence.Category, deref(evidence.ActualEgressIPv4), deref(evidence.HysteriaVersion), deref(evidence.WireGuardVersion)); err != nil { - return nil, err + return nil, &operationStageError{Stage: "observed-state", StateChanged: "route-deployed-local-state-committed", Retry: "audit", Err: err} } result := map[string]any{"route": route.ID, "state": "deployed", "enabled": true, "validation": evidence.Sanitized()} if renderClients { diff --git a/version.txt b/version.txt index c043eea..b1b25a5 100644 --- a/version.txt +++ b/version.txt @@ -1 +1 @@ -2.2.1 +2.2.2