test: basic replicated deployment, prioritise host with the most up-to-date containers

This commit is contained in:
Pavel Sviderski
2025-03-02 10:23:49 +10:00
parent a16b2d1c74
commit 0896bd27fb
2 changed files with 171 additions and 18 deletions
+9 -6
View File
@@ -96,7 +96,7 @@ func (s *RollingStrategy) planReplicated(
// Organise existing containers by machine.
containersOnMachine := make(map[string][]api.Container)
machineHasUpToDateContainer := make(map[string]bool)
upToDateContainersOnMachine := make(map[string]int)
if svc != nil {
runningSpecs := make(map[string]api.ServiceSpec)
for _, c := range svc.Containers {
@@ -108,7 +108,7 @@ func (s *RollingStrategy) planReplicated(
if err == nil {
runningSpecs[c.Container.ID] = cs
if cs.Equals(spec) {
machineHasUpToDateContainer[c.MachineID] = true
upToDateContainersOnMachine[c.MachineID] += 1
}
}
}
@@ -128,13 +128,16 @@ func (s *RollingStrategy) planReplicated(
containersOnMachine[c.MachineID] = append(containersOnMachine[c.MachineID], c.Container)
}
// Sort machines such that machines with containers that match the desired spec are first, followed by machines
// with existing containers, and finally machines without containers.
// Sort machines such that machines with the most up-to-date containers are first, followed by machines with
// existing containers, and finally machines without containers.
slices.SortFunc(availableMachines, func(m1, m2 *pb.MachineInfo) int {
if machineHasUpToDateContainer[m1.Id] {
if upToDateContainersOnMachine[m1.Id] > 0 && upToDateContainersOnMachine[m2.Id] > 0 {
return upToDateContainersOnMachine[m2.Id] - upToDateContainersOnMachine[m1.Id]
}
if upToDateContainersOnMachine[m1.Id] > 0 {
return -1
}
if machineHasUpToDateContainer[m2.Id] {
if upToDateContainersOnMachine[m2.Id] > 0 {
return 1
}
return len(containersOnMachine[m2.Id]) - len(containersOnMachine[m1.Id])
+162 -12
View File
@@ -42,7 +42,7 @@ func TestDeployment(t *testing.T) {
name := "global-deployment"
t.Cleanup(func() {
err := cli.RemoveService(ctx, name)
if errors.Is(err, client.ErrNotFound) {
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
@@ -188,7 +188,7 @@ func TestDeployment(t *testing.T) {
name := "global-deployment-filtered"
t.Cleanup(func() {
err := cli.RemoveService(ctx, name)
if errors.Is(err, client.ErrNotFound) {
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
})
@@ -303,7 +303,7 @@ func TestDeployment(t *testing.T) {
t.Run("caddy", func(t *testing.T) {
t.Cleanup(func() {
err := cli.RemoveService(ctx, client.CaddyServiceName)
if errors.Is(err, client.ErrNotFound) {
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
})
@@ -350,7 +350,7 @@ func TestDeployment(t *testing.T) {
t.Run("caddy with machine filter", func(t *testing.T) {
t.Cleanup(func() {
err := cli.RemoveService(ctx, client.CaddyServiceName)
if errors.Is(err, client.ErrNotFound) {
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
})
@@ -404,6 +404,162 @@ func TestDeployment(t *testing.T) {
assert.Equal(t, c.Machines[2].Name, machine2.Machine.Name)
})
t.Run("replicated", func(t *testing.T) {
t.Parallel()
name := "replicated-deployment"
t.Cleanup(func() {
err := cli.RemoveService(ctx, name)
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
})
// 1. Create a basic replicated service with 2 replicas.
spec := api.ServiceSpec{
Name: name,
Mode: api.ServiceModeReplicated,
Container: api.ContainerSpec{
Image: "portainer/pause:latest",
},
Replicas: 2,
}
deploy, err := cli.NewDeployment(spec, nil)
require.NoError(t, err)
err = deploy.Validate(ctx)
require.NoError(t, err)
plan, err := deploy.Plan(ctx)
require.NoError(t, err)
assert.Len(t, plan.SequenceOperation.Operations, 2) // 2 run operations for 2 replicas
svcID, err := deploy.Run(ctx)
require.NoError(t, err)
assert.NotEmpty(t, svcID)
// Verify service was created with correct settings.
svc, err := cli.InspectService(ctx, name)
require.NoError(t, err)
assert.Equal(t, name, svc.Name)
assert.Equal(t, api.ServiceModeReplicated, svc.Mode)
assert.Len(t, svc.Containers, 2, "expected 2 replicas")
// Verify containers are on different machines for balanced distribution.
machines := make(map[string]struct{})
for _, ctr := range svc.Containers {
machines[ctr.MachineID] = struct{}{}
// Verify container spec matches our deployment spec.
svcSpec, err := ctr.Container.ServiceSpec()
require.NoError(t, err)
assert.True(t, svcSpec.Equals(spec))
}
assert.Len(t, machines, 2, "containers should be on different machines")
// Store the initial container IDs.
initialContainers := make(map[string]string) // machineID -> containerID
for _, ctr := range svc.Containers {
initialContainers[ctr.MachineID] = ctr.Container.ID
}
// 2. Update the service with a new configuration.
init := true
updatedSpec := spec
updatedSpec.Container.Init = &init
deploy, err = cli.NewDeployment(updatedSpec, nil)
require.NoError(t, err)
plan, err = deploy.Plan(ctx)
require.NoError(t, err)
assert.Len(t, plan.Operations, 4, "expected 2 run + 2 remove operations")
_, err = deploy.Run(ctx)
require.NoError(t, err)
// Verify service was updated.
svc, err = cli.InspectService(ctx, name)
require.NoError(t, err)
assert.Equal(t, name, svc.Name)
assert.Len(t, svc.Containers, 2)
// Verify initial containers were updated.
for _, ctr := range svc.Containers {
initialCtr, ok := initialContainers[ctr.MachineID]
require.True(t, ok, "Updated container should have replaced one of the initial containers")
assert.NotEqual(t, initialCtr, ctr.Container.ID,
"Container on machine %s should have been updated", ctr.MachineID)
svcSpec, err := ctr.Container.ServiceSpec()
require.NoError(t, err)
assert.True(t, svcSpec.Equals(updatedSpec))
}
// 3. Update to 4 replicas with a different configuration.
initialContainers = make(map[string]string) // Reset container tracking.
for _, ctr := range svc.Containers {
initialContainers[ctr.MachineID] = ctr.Container.ID
}
fourReplicaSpec := updatedSpec
fourReplicaSpec.Container.Command = []string{"updated"}
fourReplicaSpec.Replicas = 4
deploy, err = cli.NewDeployment(fourReplicaSpec, nil)
require.NoError(t, err)
plan, err = deploy.Plan(ctx)
require.NoError(t, err)
assert.Len(t, plan.Operations, 6, "Expected 4 run + 2 remove operations")
_, err = deploy.Run(ctx)
require.NoError(t, err)
// Verify service now has 4 containers.
svc, err = cli.InspectService(ctx, name)
require.NoError(t, err)
assert.Equal(t, name, svc.Name)
assert.Len(t, svc.Containers, 4, "Expected 4 replicas")
// Count containers per machine.
machineContainerCount := make(map[string]int)
for _, ctr := range svc.Containers {
machineContainerCount[ctr.MachineID]++
// Verify all containers match the new spec
svcSpec, err := ctr.Container.ServiceSpec()
require.NoError(t, err)
assert.True(t, svcSpec.Equals(fourReplicaSpec))
// For existing machines, verify containers were replaced
if initialID, ok := initialContainers[ctr.MachineID]; ok {
assert.NotEqual(t, initialID, ctr.Container.ID,
"Container on machine %s should have been updated", ctr.MachineID)
}
}
// Verify even distributions across machines.
assert.Len(t, machineContainerCount, 3, "Expected containers on all 3 machines")
for _, count := range machineContainerCount {
assert.GreaterOrEqual(t, count, 1, "Expected at least 1 container on each machine")
}
// 4. Redeploy the exact same spec and verify it's a noop.
deploy, err = cli.NewDeployment(fourReplicaSpec, nil)
require.NoError(t, err)
plan, err = deploy.Plan(ctx)
require.NoError(t, err)
svc, err = cli.InspectService(ctx, name)
require.NoError(t, err)
assert.Empty(t, plan.Operations, "Redeploying the same spec should be a no-op")
})
// TODO: test deployments with unreachable machines. See https://github.com/psviderski/uncloud/issues/29.
}
@@ -646,12 +802,9 @@ func TestRunService(t *testing.T) {
name := "1-replica-ports"
t.Cleanup(func() {
err := cli.RemoveService(ctx, name)
if errors.Is(err, client.ErrNotFound) {
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
_, err = cli.InspectService(ctx, name)
require.ErrorIs(t, err, client.ErrNotFound)
})
spec := api.ServiceSpec{
@@ -700,12 +853,9 @@ func TestRunService(t *testing.T) {
name := "global"
t.Cleanup(func() {
err := cli.RemoveService(ctx, name)
if errors.Is(err, client.ErrNotFound) {
if !errors.Is(err, client.ErrNotFound) {
require.NoError(t, err)
}
_, err = cli.InspectService(ctx, name)
require.ErrorIs(t, err, client.ErrNotFound)
})
resp, err := cli.RunService(ctx, api.ServiceSpec{