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
82 changes: 82 additions & 0 deletions internal/steward/deployment_recovery_test.go
Original file line number Diff line number Diff line change
@@ -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))
}
}
46 changes: 37 additions & 9 deletions internal/steward/execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion version.txt
Original file line number Diff line number Diff line change
@@ -1 +1 @@
2.2.1
2.2.2