reconcile: fix scale-down propagation, surface hash errors, dedupe stops

Three fixes surfaced by an external code review of the new reconciler:

1. Scale-down now propagates through serviceNodes. When a service is
   scaled down (all containers in excess), reconcileService used to
   continue without assigning lastNode, leaving r.serviceNodes[svc]
   unset. Dependent services then declared no edge on the scale-down
   ops and could start before the cleanup finished. Track the
   RemoveContainer node as lastNode so depends_on chains pick it up.

2. mustRecreate errors are no longer silently ignored by sortContainers.
   The comparator used `obsi, _ := r.mustRecreate(...)`, falling back
   to false on any hashing error. Pre-compute obsolescence into a map
   keyed by container ID before sorting and propagate the error to
   reconcileService.

3. A container is no longer Stopped twice when its network and its
   config both diverge. planRecreateNetwork already stops the affected
   container as part of the disconnect/remove/recreate dance; the
   subsequent planRecreateContainer (triggered via hasNetworkMismatch)
   used to add another OpStopContainer against the now-stopped target.
   Track stops in r.stoppedByPlan; planRecreateContainer reuses an
   existing Stop node when present, and chains its Remove on both that
   Stop and the replacement Create.

Two golden tests (TestReconcileNetworks_Diverged*) are updated to
reflect the new, dedupe'd plan shape (one Stop instead of two per
recreated container).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Signed-off-by: Guillaume Lours <glours@users.noreply.github.com>
This commit is contained in:
Guillaume Lours 2026-05-28 16:59:55 +02:00 committed by Guillaume Lours
parent c96ee45f51
commit a06368c333
2 changed files with 75 additions and 36 deletions

View file

@ -85,6 +85,11 @@ type reconciler struct {
// serviceNodes tracks the last plan node per service, so dependent
// services can order their operations after dependencies.
serviceNodes map[string]*PlanNode
// stoppedByPlan records containers already stopped by an earlier stage
// of the plan (typically planRecreateNetwork) so that downstream stages
// can chain on the existing OpStopContainer instead of emitting a second
// one against an already-stopped container.
stoppedByPlan map[string]*PlanNode // container ID → existing Stop node
}
// reconcile is the main entry point: it builds a Plan from desired vs observed state.
@ -92,14 +97,15 @@ type reconciler struct {
// reconciler.prompt field).
func reconcile(_ context.Context, project *types.Project, observed *ObservedState, options ReconcileOptions, prompt Prompt) (*Plan, error) {
r := &reconciler{
project: project,
observed: observed,
options: options,
prompt: prompt,
plan: &Plan{},
networkNodes: map[string]*PlanNode{},
volumeNodes: map[string]*PlanNode{},
serviceNodes: map[string]*PlanNode{},
project: project,
observed: observed,
options: options,
prompt: prompt,
plan: &Plan{},
networkNodes: map[string]*PlanNode{},
volumeNodes: map[string]*PlanNode{},
serviceNodes: map[string]*PlanNode{},
stoppedByPlan: map[string]*PlanNode{},
}
if err := r.reconcileNetworks(); err != nil {
@ -166,7 +172,9 @@ func (r *reconciler) planRecreateNetwork(key string, nw *types.NetworkConfig) er
affectedServices := r.servicesUsingNetwork(key)
affectedContainers := r.containersForServices(affectedServices)
// Stop all affected containers
// Stop all affected containers, recording each Stop node so that a later
// recreate of the same container does not emit a second Stop against a
// container that is already stopped.
var stopNodes []*PlanNode
for i := range affectedContainers {
oc := &affectedContainers[i]
@ -177,6 +185,7 @@ func (r *reconciler) planRecreateNetwork(key string, nw *types.NetworkConfig) er
Container: &oc.Summary,
}, "")
stopNodes = append(stopNodes, node)
r.stoppedByPlan[oc.ID] = node
}
// Disconnect all affected containers from the *observed* network (each depends on its own stop)
@ -427,7 +436,9 @@ func (r *reconciler) reconcileService(service types.ServiceConfig) error {
// Sort containers: obsolete first, then by number descending, then reverse
// to get the same ordering as the existing convergence code.
r.sortContainers(containers, service, strategy)
if err := r.sortContainers(containers, service, strategy); err != nil {
return err
}
// Collect dependency nodes that container creation should depend on
infraDeps := r.infrastructureDeps(service)
@ -437,7 +448,9 @@ func (r *reconciler) reconcileService(service types.ServiceConfig) error {
// Process existing containers
for i, oc := range containers {
if i >= expected {
// Scale down: stop + remove excess containers
// Scale down: stop + remove excess containers. Track the remove
// node so dependent services wait for the scale-down to finish
// even when no other operation runs on this service.
stopNode := r.plan.addNode(Operation{
Type: OpStopContainer,
ResourceID: fmt.Sprintf("service:%s:%d", service.Name, oc.Number),
@ -445,7 +458,7 @@ func (r *reconciler) reconcileService(service types.ServiceConfig) error {
Container: &containers[i].Summary,
Timeout: r.options.Timeout,
}, "")
r.plan.addNode(Operation{
lastNode = r.plan.addNode(Operation{
Type: OpRemoveContainer,
ResourceID: fmt.Sprintf("service:%s:%d", service.Name, oc.Number),
Cause: "scale down",
@ -610,22 +623,36 @@ func (r *reconciler) planRecreateContainer(service types.ServiceConfig, oc *Obse
Name: tmpName,
}, group, allDeps...)
// 2. Stop old container
stopNode := r.plan.addNode(Operation{
Type: OpStopContainer,
ResourceID: resID,
Cause: fmt.Sprintf("replaced by #%d", createNode.ID),
Container: &oc.Summary,
Timeout: r.options.Timeout,
}, group, createNode)
// 2. Stop old container. If an earlier stage of the plan (e.g.
// planRecreateNetwork) already scheduled a Stop for this container,
// reuse it instead of emitting a second one against an already-stopped
// container. The reused node carries no group, which is fine: the
// recreate's group tracker still drives Working/Done from the create.
stopNode, alreadyStopped := r.stoppedByPlan[oc.ID]
if !alreadyStopped {
stopNode = r.plan.addNode(Operation{
Type: OpStopContainer,
ResourceID: resID,
Cause: fmt.Sprintf("replaced by #%d", createNode.ID),
Container: &oc.Summary,
Timeout: r.options.Timeout,
}, group, createNode)
r.stoppedByPlan[oc.ID] = stopNode
}
// 3. Remove old container
// 3. Remove old container. Depend on the (existing or new) stop and on
// the create node so the new container is in place before the old one
// goes away.
removeDeps := []*PlanNode{stopNode}
if alreadyStopped {
removeDeps = append(removeDeps, createNode)
}
removeNode := r.plan.addNode(Operation{
Type: OpRemoveContainer,
ResourceID: resID,
Cause: fmt.Sprintf("replaced by #%d", createNode.ID),
Container: &oc.Summary,
}, group, stopNode)
}, group, removeDeps...)
// 4. Rename to final name. Link to the create node so the executor can
// fetch the resulting container ID directly.
@ -691,10 +718,21 @@ func (r *reconciler) infrastructureDeps(service types.ServiceConfig) []*PlanNode
// sortContainers sorts containers the same way as convergence.go:138-160:
// obsolete first, then by container number descending, then reversed.
func (r *reconciler) sortContainers(containers []ObservedContainer, service types.ServiceConfig, policy string) {
//
// mustRecreate is evaluated once per container before sorting, both to avoid
// quadratic re-evaluation in the comparator and to surface any hashing error
// instead of silently treating the container as non-obsolete.
func (r *reconciler) sortContainers(containers []ObservedContainer, service types.ServiceConfig, policy string) error {
obsolete := make(map[string]bool, len(containers))
for _, oc := range containers {
o, err := r.mustRecreate(service, oc, policy)
if err != nil {
return err
}
obsolete[oc.ID] = o
}
sort.Slice(containers, func(i, j int) bool {
obsi, _ := r.mustRecreate(service, containers[i], policy)
obsj, _ := r.mustRecreate(service, containers[j], policy)
obsi, obsj := obsolete[containers[i].ID], obsolete[containers[j].ID]
if obsi != obsj {
return obsi // obsolete first
}
@ -705,6 +743,7 @@ func (r *reconciler) sortContainers(containers []ObservedContainer, service type
return containers[i].Summary.Created < containers[j].Summary.Created
})
slices.Reverse(containers)
return nil
}
// reconcileOrphans plans stop + remove for orphaned containers.

View file

@ -146,15 +146,16 @@ func TestReconcileNetworks_Diverged(t *testing.T) {
plan, err := reconcile(t.Context(), project, observed, defaultReconcileOptions(), noPrompt)
assert.NilError(t, err)
// The recreate phase reuses the Stop from the network-recreate phase
// instead of emitting a second one against an already-stopped container.
assert.Equal(t, plan.String(), strings.TrimSpace(`
[] -> #1 service:web:1, StopContainer, network frontend config changed
[1] -> #2 service:web:1, DisconnectNetwork, network frontend recreate
[2] -> #3 network:frontend, RemoveNetwork, config hash diverged
[3] -> #4 network:frontend, CreateNetwork, recreate after config change
[4] -> #5 service:web:1, CreateContainer, config changed (tmpName) [recreate:web:1]
[5] -> #6 service:web:1, StopContainer, replaced by #5 [recreate:web:1]
[6] -> #7 service:web:1, RemoveContainer, replaced by #5 [recreate:web:1]
[7] -> #8 service:web:1, RenameContainer, finalize recreate [recreate:web:1]
[1,5] -> #6 service:web:1, RemoveContainer, replaced by #5 [recreate:web:1]
[6] -> #7 service:web:1, RenameContainer, finalize recreate [recreate:web:1]
`)+"\n")
}
@ -198,7 +199,8 @@ func TestReconcileNetworks_DivergedMultipleServices(t *testing.T) {
plan, err := reconcile(t.Context(), project, observed, defaultReconcileOptions(), noPrompt)
assert.NilError(t, err)
// Services sorted alphabetically: api before web
// Services sorted alphabetically: api before web. Each service's recreate
// reuses the Stop from the network-recreate phase (no second Stop).
assert.Equal(t, plan.String(), strings.TrimSpace(`
[] -> #1 service:api:1, StopContainer, network frontend config changed
[] -> #2 service:web:1, StopContainer, network frontend config changed
@ -207,13 +209,11 @@ func TestReconcileNetworks_DivergedMultipleServices(t *testing.T) {
[3,4] -> #5 network:frontend, RemoveNetwork, config hash diverged
[5] -> #6 network:frontend, CreateNetwork, recreate after config change
[6] -> #7 service:api:1, CreateContainer, config changed (tmpName) [recreate:api:1]
[7] -> #8 service:api:1, StopContainer, replaced by #7 [recreate:api:1]
[8] -> #9 service:api:1, RemoveContainer, replaced by #7 [recreate:api:1]
[9] -> #10 service:api:1, RenameContainer, finalize recreate [recreate:api:1]
[6] -> #11 service:web:1, CreateContainer, config changed (tmpName) [recreate:web:1]
[11] -> #12 service:web:1, StopContainer, replaced by #11 [recreate:web:1]
[12] -> #13 service:web:1, RemoveContainer, replaced by #11 [recreate:web:1]
[13] -> #14 service:web:1, RenameContainer, finalize recreate [recreate:web:1]
[1,7] -> #8 service:api:1, RemoveContainer, replaced by #7 [recreate:api:1]
[8] -> #9 service:api:1, RenameContainer, finalize recreate [recreate:api:1]
[6] -> #10 service:web:1, CreateContainer, config changed (tmpName) [recreate:web:1]
[2,10] -> #11 service:web:1, RemoveContainer, replaced by #10 [recreate:web:1]
[11] -> #12 service:web:1, RenameContainer, finalize recreate [recreate:web:1]
`)+"\n")
}