From d28d74d6a49632a850c87d30c0e98570510ca36b Mon Sep 17 00:00:00 2001 From: Nicolas De Loof Date: Wed, 30 Sep 2026 12:49:27 +0200 Subject: [PATCH 1/5] feat: executor runs start-phase operations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wire the plan executor to the start-phase vocabulary the reconciler already produces (#14200): OpWaitCondition, OpRunPreStart, the enriched OpStartContainer, and OpRunPostStart all had no case in executeNode and hit "unknown operation type". Every operation reuses the exact imperative primitive it mirrors — waitDependency, runPreStart, runHook, injectSecrets/injectConfigs — so behavior and events cannot drift between the two engines: execWaitCondition resolves the dependency's live containers and hands off to the same polling/condition-check functions waitDependencies uses; the enriched OpStartContainer resolves its target either from the observed container or, for a replica the plan itself creates, from the CreateContainer node's result (the mechanism OpRenameContainer already relies on), then injects secrets/configs before ContainerStart; OpRunPreStart re-checks the daemon first so a replica that started in the observe-to-execute window skips a redundant hook run. A replica's pre_start/start/post_start chain shares one event group, reporting a single Starting→Started progression exactly like startServiceContainer does today — Starting fires specifically at StartContainer, not at the (silent) pre_start hooks preceding it. pkg/api/event.go documents this per-resource progression contract. Still inert in production: create() (executePlan's only caller never builds a plan with start-phase scope, so none of this runs outside tests until a later brick in epic #14081 wires a real caller. EOF ) Signed-off-by: Nicolas De Loof --- pkg/api/event.go | 22 ++ pkg/compose/executor.go | 26 +- pkg/compose/executor_events.go | 114 ++++++-- pkg/compose/executor_ops.go | 144 +++++++++- pkg/compose/executor_start_test.go | 437 +++++++++++++++++++++++++++++ pkg/compose/executor_test.go | 10 +- 6 files changed, 725 insertions(+), 28 deletions(-) create mode 100644 pkg/compose/executor_start_test.go diff --git a/pkg/api/event.go b/pkg/api/event.go index 7ccb1138210..c300a4cc95c 100644 --- a/pkg/api/event.go +++ b/pkg/api/event.go @@ -97,6 +97,28 @@ func (e *Resource) StatusText() string { } // EventProcessor is notified about Compose operations and tasks +// +// # Event contract +// +// The event stream is ordered per resource, not globally: operations run +// concurrently and events of distinct resources interleave freely, but a +// given resource's events follow a fixed progression. Consumers must key on +// Resource.ID and must not rely on global ordering. +// +// Typical per-resource progressions (statuses from the Status* constants): +// +// container being created: Creating → Created +// container being started: Starting → Started (pre_start/post_start +// hooks, when declared, run silently within it) +// container being replaced: Recreate → Recreated +// container stop/removal: Stopping → Stopped, Removing → Removed +// dependency being waited: Waiting → Healthy | Exited +// networks and volumes: Creating → Created, Removing → Removed +// +// Any progression can end early with an Error status, or — for optional +// outcomes such as a dependency declared with required: false — with a +// Skipped status carrying the reason. Status texts are stable identifiers: +// tools may match on them. type EventProcessor interface { // Start is triggered as a Compose operation is starting with context Start(ctx context.Context, operation string) diff --git a/pkg/compose/executor.go b/pkg/compose/executor.go index f69cd155c4b..a59ca9e9090 100644 --- a/pkg/compose/executor.go +++ b/pkg/compose/executor.go @@ -23,6 +23,8 @@ import ( "github.com/compose-spec/compose-go/v2/types" "golang.org/x/sync/errgroup" + + "github.com/docker/compose/v5/pkg/api" ) // planExecutor executes a reconciliation Plan by walking the DAG and performing @@ -33,6 +35,13 @@ type planExecutor struct { project *types.Project pctx *reconciliationContext + // listener streams pre_start/post_start hook logs, exactly like the + // imperative start path. nil until an attached caller wires one in (the + // interactive-up convergence, a later lot of #14081) — every caller today + // passes nil, so hook execution stays silent, matching today's detached + // behavior. + listener api.ContainerEventListener + // containersByService is a live view used to resolve service references // (network_mode: service:x, volumes_from, ipc, pid) without a daemon // round-trip per create. @@ -68,17 +77,24 @@ func (pc *reconciliationContext) get(nodeID int) operationResult { // executePlan walks the plan DAG, executing nodes in parallel where possible // while respecting dependency edges. It emits progress events and handles // group-based event aggregation for composite operations like recreate. +// +// No caller threads a hook-log listener through here yet: create() (its only +// caller) plans no start-phase operations today. newPlanExecutor takes one +// directly for that reason — it is what the interactive-up convergence (a +// later lot of #14081) will call once it needs to stream pre_start/post_start +// hook logs into an attached session. func (s *composeService) executePlan(ctx context.Context, project *types.Project, observed *ObservedState, plan *Plan) error { - return s.newPlanExecutor(project, observed).run(ctx, plan) + return s.newPlanExecutor(project, observed, nil).run(ctx, plan) } // newPlanExecutor constructs a planExecutor seeded from the observed state. // Split out from executePlan so tests can inspect the executor's live state // (e.g. the containersByService cache) after running a plan. -func (s *composeService) newPlanExecutor(project *types.Project, observed *ObservedState) *planExecutor { +func (s *composeService) newPlanExecutor(project *types.Project, observed *ObservedState, listener api.ContainerEventListener) *planExecutor { return &planExecutor{ compose: s, project: project, + listener: listener, pctx: &reconciliationContext{results: map[int]operationResult{}}, containersByService: observed.containersByService(), } @@ -174,6 +190,12 @@ func (exec *planExecutor) executeNode(ctx context.Context, node *PlanNode) error return exec.execRenameContainer(ctx, node) case OpCreateHookContainer: return exec.execCreateHookContainer(ctx, node) + case OpWaitCondition: + return exec.execWaitCondition(ctx, op) + case OpRunPreStart: + return exec.execRunPreStart(ctx, op) + case OpRunPostStart: + return exec.execRunPostStart(ctx, op) case OpRunProvider: return exec.compose.runPlugin(ctx, exec.project, *op.Service, "up") default: diff --git a/pkg/compose/executor_events.go b/pkg/compose/executor_events.go index 5177d09313e..cf846c57f2e 100644 --- a/pkg/compose/executor_events.go +++ b/pkg/compose/executor_events.go @@ -22,44 +22,96 @@ import ( "github.com/docker/compose/v5/pkg/api" ) -// groupTracker manages event emission for grouped nodes (e.g. recreate). -// The first node starting emits Working, the last finishing emits Done. +// groupTracker manages event emission for grouped nodes (e.g. recreate, or a +// replica's start chain). The first node starting emits Working, the last +// finishing emits Done. type groupTracker struct { mu sync.Mutex groups map[string]*groupState + exec *planExecutor // resolves the event name of a not-yet-materialized container +} + +// groupKind picks the Working/Done status text for a group: distinct +// composite operations get distinct progressions on the event contract +// (recreate: "Recreate"/"Recreated"; a replica's start chain — pre_start, +// start, post_start folded into one line — the same Starting/Started +// progression a plain start reports). +type groupKind int + +const ( + groupRecreate groupKind = iota + groupStart +) + +func groupKindOf(t OperationType) groupKind { + switch t { + case OpRunPreStart, OpStartContainer, OpRunPostStart: + return groupStart + default: + return groupRecreate + } +} + +func (k groupKind) workingText() string { + if k == groupStart { + return api.StatusStarting + } + return "Recreate" +} + +func (k groupKind) doneText() string { + if k == groupStart { + return api.StatusStarted + } + return "Recreated" } type groupState struct { - eventName string // e.g. "Container myproject-web-1" + kind groupKind + eventName string // e.g. "Container myproject-web-1"; resolved lazily for a start group (see groupEventName) total int // total nodes in this group started int // nodes that have started done int // nodes that have completed } func (exec *planExecutor) buildGroupTracker(plan *Plan) *groupTracker { - gt := &groupTracker{groups: map[string]*groupState{}} + gt := &groupTracker{groups: map[string]*groupState{}, exec: exec} for _, node := range plan.Nodes { if node.Group == "" { continue } - if _, ok := gt.groups[node.Group]; !ok { - gt.groups[node.Group] = &groupState{} + gs, ok := gt.groups[node.Group] + if !ok { + gs = &groupState{kind: groupKindOf(node.Operation.Type)} + gt.groups[node.Group] = gs } - gt.groups[node.Group].total++ - // Pick the event name from a node that has the existing container reference - if gt.groups[node.Group].eventName == "" && node.Operation.Container != nil { - gt.groups[node.Group].eventName = getContainerProgressName(*node.Operation.Container) - } - } - // Fallback for groups where no node had a Container (shouldn't happen for recreate) - for name, gs := range gt.groups { - if gs.eventName == "" { - gs.eventName = name + gs.total++ + // Pick the event name from a node that has the existing container + // reference. A start group's first node may still be a plan-created + // replica with no Summary yet — its name resolves lazily, once + // execution reaches it (see onNodeStart). + if gs.eventName == "" && node.Operation.Container != nil { + gs.eventName = getContainerProgressName(*node.Operation.Container) } } return gt } +// groupEventName names a group from its first executing node, for the case +// buildGroupTracker could not resolve statically: a start-phase node +// targeting a replica the plan itself creates. By the time this node starts, +// the CreateContainer node it depends on has already run and published its +// result (see resolveContainerID) — the DAG dependency guarantees it. +func (exec *planExecutor) groupEventName(op Operation) string { + if op.Container != nil { + return getContainerProgressName(*op.Container) + } + if name := exec.pctx.get(op.CreateNodeID).ContainerName; name != "" { + return "Container " + name + } + return op.ResourceID +} + func (gt *groupTracker) onNodeStart(node *PlanNode, events api.EventProcessor) { if node.Group == "" { // Ungrouped: emit individual event @@ -69,10 +121,26 @@ func (gt *groupTracker) onNodeStart(node *PlanNode, events api.EventProcessor) { gt.mu.Lock() defer gt.mu.Unlock() gs := gt.groups[node.Group] + if gs.eventName == "" { + gs.eventName = gt.exec.groupEventName(node.Operation) + } gs.started++ - if gs.started == 1 { - events.On(newEvent(gs.eventName, api.Working, "Recreate")) + if gs.triggersWorking(node.Operation.Type) { + events.On(newEvent(gs.eventName, api.Working, gs.kind.workingText())) + } +} + +// triggersWorking reports whether this node starting should fire the group's +// Working event. A recreate group fires on its first node; a start group +// fires specifically on OpStartContainer — pre_start hooks, when planned, +// run silently before it, matching the imperative engine's startService, +// which never surfaces a Starting event until the ContainerStart call itself +// begins. +func (gs *groupState) triggersWorking(t OperationType) bool { + if gs.kind == groupStart { + return t == OpStartContainer } + return gs.started == 1 } func (gt *groupTracker) onNodeDone(node *PlanNode, events api.EventProcessor) { @@ -85,7 +153,7 @@ func (gt *groupTracker) onNodeDone(node *PlanNode, events api.EventProcessor) { gs := gt.groups[node.Group] gs.done++ if gs.done == gs.total { - events.On(newEvent(gs.eventName, api.Done, "Recreated")) + events.On(newEvent(gs.eventName, api.Done, gs.kind.doneText())) } } @@ -155,6 +223,14 @@ func emitDoneEvent(node *PlanNode, events api.EventProcessor) { // emitErrorEvent emits an error event for an ungrouped node. func emitErrorEvent(node *PlanNode, events api.EventProcessor, err error) { op := node.Operation + if op.Type == OpWaitCondition { + // execWaitCondition already reported this failure as one event per + // container of the dependency it was waiting on (see + // waitDependency/checkDependency*), exactly like waitDependencies + // does today. A second, generic event on "wait:..." here would be a + // confusing duplicate with no matching resource. + return + } var id string switch { case op.Container != nil: diff --git a/pkg/compose/executor_ops.go b/pkg/compose/executor_ops.go index d9f61c8a660..9b303109110 100644 --- a/pkg/compose/executor_ops.go +++ b/pkg/compose/executor_ops.go @@ -21,6 +21,7 @@ import ( "fmt" "slices" + "github.com/compose-spec/compose-go/v2/types" "github.com/containerd/errdefs" "github.com/moby/moby/api/types/container" "github.com/moby/moby/client" @@ -127,13 +128,152 @@ func (exec *planExecutor) execCreateContainer(ctx context.Context, node *PlanNod return nil } +// execStartContainer starts a container. A bare operation (no Service — +// today only the create phase's exceptional-state restart of a paused/dead +// container) is a plain ContainerStart. A start-phase operation (Service set) +// performs the full service start — secret/config injection right before +// ContainerStart — mirroring startServiceContainer; its target resolves +// either from the observed container or, for a replica the plan itself +// creates, from the CreateContainer node's result (the same mechanism +// OpRenameContainer already uses). func (exec *planExecutor) execStartContainer(ctx context.Context, op Operation) error { + if op.Service == nil { + startMx.Lock() + defer startMx.Unlock() + _, err := exec.compose.apiClient().ContainerStart(ctx, op.Container.ID, client.ContainerStartOptions{}) + return err + } + + id, err := exec.resolveContainerID(op) + if err != nil { + return err + } + if err := exec.compose.injectSecrets(ctx, exec.project, *op.Service, id); err != nil { + return err + } + if err := exec.compose.injectConfigs(ctx, exec.project, *op.Service, id); err != nil { + return err + } + startMx.Lock() - defer startMx.Unlock() - _, err := exec.compose.apiClient().ContainerStart(ctx, op.Container.ID, client.ContainerStartOptions{}) + _, err = exec.compose.apiClient().ContainerStart(ctx, id, client.ContainerStartOptions{}) + startMx.Unlock() return err } +// resolveContainerID returns the ID of the container an operation targets: +// the observed container when the reconciler already had one, otherwise the +// result of the CreateContainer node it references. +func (exec *planExecutor) resolveContainerID(op Operation) (string, error) { + if op.Container != nil { + return op.Container.ID, nil + } + res := exec.pctx.get(op.CreateNodeID) + if res.ContainerID == "" { + return "", fmt.Errorf("internal: no materialized container for %s", op.ResourceID) + } + return res.ContainerID, nil +} + +// resolveContainerSummary is resolveContainerID plus the rest of the +// container.Summary that hook execution (runHook) needs — preferring the +// observed container, falling back to the live view populated by the create +// node execCreateContainer already ran (a dependency of every start-phase +// node targeting that replica). +func (exec *planExecutor) resolveContainerSummary(op Operation) (container.Summary, error) { + if op.Container != nil { + return *op.Container, nil + } + id, err := exec.resolveContainerID(op) + if err != nil { + return container.Summary{}, err + } + exec.containersMu.Lock() + defer exec.containersMu.Unlock() + for _, c := range exec.containersByService[op.Service.Name] { + if c.ID == id { + return c, nil + } + } + return container.Summary{ID: id, Names: []string{"/" + exec.pctx.get(op.CreateNodeID).ContainerName}}, nil +} + +// execWaitCondition polls the dependency service named by the operation +// until it satisfies the declared depends_on condition or the context ends — +// the plan-side equivalent of one waitDependencies edge, reusing the exact +// same polling primitive (waitDependency) and per-condition checks the +// imperative engine uses, so events and error messages cannot drift between +// the two. required: false dependencies are marked BestEffort by the +// reconciler: a missing dependency or a failed/timed-out condition is then a +// warning, not a plan failure — matching waitDependencies' own +// optional-dependency handling. +// +// Unlike waitDependencies, this applies no deadline of its own: nothing +// produces one yet (no ReconcileOptions field feeds a per-wait timeout the +// way api.CreateOptions.WaitTimeout does today). A future caller needing that +// — e.g. `up --wait` once it runs on the plan — wraps ctx before executing +// the plan, or adds a Timeout to the operation for execWaitCondition to wrap +// here. +func (exec *planExecutor) execWaitCondition(ctx context.Context, op Operation) error { + s := exec.compose + exec.containersMu.Lock() + waitingFor := exec.containersByService[op.Name].filter(isNotOneOff, isNotHookContainer) + exec.containersMu.Unlock() + + config := types.ServiceDependency{Condition: op.Condition, Required: !op.BestEffort} + + if len(waitingFor) == 0 { + if config.Required { + return fmt.Errorf("missing dependency %s", op.Name) + } + logrus.Warnf("missing dependency %s", op.Name) + return nil + } + + s.events.On(containerEvents(waitingFor, waiting)...) + // The wait node is shared across every dependent awaiting the same + // (dependency, condition) pair (see waitConditionNode), so no single + // requester name would be accurate; op.ResourceID identifies the wait + // itself instead. waitDependency only ever reads this for one + // practically unreachable log line (an unsupported depends_on condition, + // filtered out before a plan is ever built). + return s.waitDependency(ctx, op.ResourceID, op.Name, config, waitingFor) +} + +// execRunPreStart runs the service's pre_start hooks against the runner +// containers the create phase prepared (see execCreateHookContainer). The +// "once per service, no replica running at observation" rule is a plan-time +// decision; this re-checks the daemon first so a replica started in the +// window between observation and execution — the drift the plan design +// accepts — skips a second, redundant run of the hooks rather than erroring +// on runners already consumed. +func (exec *planExecutor) execRunPreStart(ctx context.Context, op Operation) error { + running, err := exec.compose.getContainers(ctx, exec.project.Name, oneOffExclude, false, op.Service.Name) + if err != nil { + return err + } + if len(running) > 0 { + logrus.Debugf("skipping pre_start hooks of service %s: a replica is already running", op.Service.Name) + return nil + } + return exec.compose.runPreStart(ctx, exec.project, *op.Service, exec.listener) +} + +// execRunPostStart runs the service's post_start hooks against the replica +// the start chain just brought up. +func (exec *planExecutor) execRunPostStart(ctx context.Context, op Operation) error { + ctr, err := exec.resolveContainerSummary(op) + if err != nil { + return err + } + for _, hook := range op.Service.PostStart { + if err := exec.compose.runHook(ctx, ctr, *op.Service, hook, exec.listener); err != nil { + return err + } + } + return nil +} + func (exec *planExecutor) execStopContainer(ctx context.Context, op Operation) error { _, err := exec.compose.apiClient().ContainerStop(ctx, op.Container.ID, client.ContainerStopOptions{ Timeout: utils.DurationSecondToInt(op.Timeout), diff --git a/pkg/compose/executor_start_test.go b/pkg/compose/executor_start_test.go new file mode 100644 index 00000000000..2e61707c241 --- /dev/null +++ b/pkg/compose/executor_start_test.go @@ -0,0 +1,437 @@ +//go:build !windows + +/* + Copyright 2020 Docker Compose CLI authors + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +*/ + +package compose + +import ( + "context" + "net" + "sync" + "testing" + + "github.com/compose-spec/compose-go/v2/types" + "github.com/docker/cli/cli/config/configfile" + "github.com/moby/moby/api/types/container" + "github.com/moby/moby/client" + "go.uber.org/mock/gomock" + "gotest.tools/v3/assert" + + "github.com/docker/compose/v5/pkg/api" + "github.com/docker/compose/v5/pkg/mocks" +) + +// recordingEvents captures, per resource ID, the sequence of status texts — +// the per-resource progression the event contract promises (see +// api.EventProcessor). +type recordingEvents struct { + mu sync.Mutex + byID map[string][]string +} + +func newRecordingEvents() *recordingEvents { return &recordingEvents{byID: map[string][]string{}} } + +func (r *recordingEvents) Start(_ context.Context, _ string) {} +func (r *recordingEvents) Done(_ string, _ bool) {} +func (r *recordingEvents) On(events ...api.Resource) { + r.mu.Lock() + defer r.mu.Unlock() + for _, e := range events { + r.byID[e.ID] = append(r.byID[e.ID], e.Text) + } +} + +func newStartPhaseTestService(t *testing.T) (*composeService, *mocks.MockAPIClient, *recordingEvents) { + t.Helper() + mockCtrl := gomock.NewController(t) + cli := mocks.NewMockCli(mockCtrl) + apiClient := mocks.NewMockAPIClient(mockCtrl) + cli.EXPECT().Client().Return(apiClient).AnyTimes() + cli.EXPECT().ConfigFile().Return(&configfile.ConfigFile{}).AnyTimes() + apiClient.EXPECT().DaemonHost().Return("unix:///var/run/docker.sock").AnyTimes() + apiClient.EXPECT().Ping(gomock.Any(), client.PingOptions{NegotiateAPIVersion: true}). + Return(client.PingResult{APIVersion: "1.44"}, nil).AnyTimes() + apiClient.EXPECT().ClientVersion().Return("1.44").AnyTimes() + + recorder := newRecordingEvents() + svcAny, err := NewComposeService(cli, WithEventProcessor(recorder)) + assert.NilError(t, err) + return svcAny.(*composeService), apiClient, recorder +} + +// TestExecWaitCondition_RequiredMissingDependencyFails verifies that a +// required depends_on condition whose dependency has no live container at +// all is a hard plan failure, exactly like waitDependencies today. +func TestExecWaitCondition_RequiredMissingDependencyFails(t *testing.T) { + svc, _, _ := newStartPhaseTestService(t) + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil) + + err := exec.execWaitCondition(t.Context(), Operation{ + Name: "db", + Condition: types.ServiceConditionHealthy, + }) + assert.ErrorContains(t, err, "missing dependency db") +} + +// TestExecWaitCondition_OptionalMissingDependencyIsTolerated mirrors +// waitDependencies' handling of an optional (required: false) dependency +// with no live container: a warning, not a plan failure, and the wait +// resolves as satisfied. +func TestExecWaitCondition_OptionalMissingDependencyIsTolerated(t *testing.T) { + svc, _, recorder := newStartPhaseTestService(t) + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil) + + err := exec.execWaitCondition(t.Context(), Operation{ + Name: "db", + Condition: types.ServiceConditionHealthy, + BestEffort: true, + }) + assert.NilError(t, err) + assert.Equal(t, len(recorder.byID), 0, "no dependency container exists to report an event against") +} + +// TestExecWaitCondition_HealthyConditionSatisfied drives execWaitCondition +// against a live dependency container that reports healthy on the first +// probe, reusing the exact same isServiceHealthy/waitDependency primitives +// the imperative engine's waitDependencies uses. +func TestExecWaitCondition_HealthyConditionSatisfied(t *testing.T) { + svc, apiClient, recorder := newStartPhaseTestService(t) + + dbSummary := container.Summary{ + ID: "db-id", + Names: []string{"/test-db-1"}, + Labels: map[string]string{api.ServiceLabel: "db", api.OneoffLabel: "False"}, + } + observed := &ObservedState{ + ProjectName: "test", + Containers: map[string][]ObservedContainer{}, + Networks: map[string][]ObservedNetwork{}, + Volumes: map[string][]ObservedVolume{}, + } + apiClient.EXPECT().ContainerInspect(gomock.Any(), "db-id", gomock.Any()).Return(client.ContainerInspectResult{ + Container: container.InspectResponse{ + ID: "db-id", + Name: "/test-db-1", + State: &container.State{ + Status: container.StateRunning, + Health: &container.Health{Status: container.Healthy}, + }, + Config: &container.Config{Healthcheck: &container.HealthConfig{Test: []string{"CMD", "true"}}}, + }, + }, nil) + + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, observed, nil) + exec.containersByService["db"] = Containers{dbSummary} + + err := exec.execWaitCondition(t.Context(), Operation{ + Name: "db", + Condition: types.ServiceConditionHealthy, + }) + assert.NilError(t, err) + assert.DeepEqual(t, recorder.byID["Container test-db-1"], []string{api.StatusWaiting, api.StatusHealthy}) +} + +// TestExecStartContainer_EnrichedResolvesCreateNodeAndStarts verifies that a +// start-phase OpStartContainer (Service set) resolves its target through the +// CreateContainer node it references — the same reconciliationContext +// mechanism OpRenameContainer already uses — rather than requiring an +// observed container.Summary. +func TestExecStartContainer_EnrichedResolvesCreateNodeAndStarts(t *testing.T) { + svc, apiClient, _ := newStartPhaseTestService(t) + + service := types.ServiceConfig{Name: "web", ContainerSpec: types.ContainerSpec{Image: "alpine"}} + project := &types.Project{Name: "test", Services: types.Services{"web": service}} + + apiClient.EXPECT().ContainerCreate(gomock.Any(), gomock.Any()). + Return(client.ContainerCreateResult{ID: "new-id"}, nil) + apiClient.EXPECT().ContainerInspect(gomock.Any(), "new-id", gomock.Any()). + Return(client.ContainerInspectResult{Container: container.InspectResponse{ + ID: "new-id", + Name: "/test-web-1", + Config: &container.Config{}, + NetworkSettings: &container.NetworkSettings{}, + }}, nil) + apiClient.EXPECT().ContainerStart(gomock.Any(), "new-id", gomock.Any()). + Return(client.ContainerStartResult{}, nil) + + plan := &Plan{} + create := plan.addNode(Operation{ + Type: OpCreateContainer, + ResourceID: "service:web:1", + Cause: "no existing container", + Service: &service, + Name: "test-web-1", + Number: 1, + }, "") + plan.addNode(Operation{ + Type: OpStartContainer, + ResourceID: "service:web:1", + Cause: "start", + Service: &service, + CreateNodeID: create.ID, + }, "start:web:1", create) + + err := svc.executePlan(t.Context(), project, emptyObservedState("test"), plan) + assert.NilError(t, err) +} + +// TestExecRunPreStart_SkipsWhenReplicaAlreadyRunning covers the +// observe-to-execute drift guard: the plan scheduled RunPreStart because no +// replica was running at observation time, but a replica started in the +// meantime — the once-per-service rule says the hooks must not run. +func TestExecRunPreStart_SkipsWhenReplicaAlreadyRunning(t *testing.T) { + svc, apiClient, _ := newStartPhaseTestService(t) + + service := types.ServiceConfig{ + Name: "web", + PreStart: []types.PreStartHook{{ContainerSpec: types.ContainerSpec{Command: types.ShellCommand{"init"}}}}, + } + running := container.Summary{ID: "c1", Labels: map[string]string{api.ServiceLabel: "web", api.OneoffLabel: "False"}} + apiClient.EXPECT().ContainerList(gomock.Any(), gomock.Any()). + Return(client.ContainerListResult{Items: []container.Summary{running}}, nil) + // No further daemon call: the hook runner scan (and any wait/start/remove) + // never happens once the drift guard skips the node. + + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil) + err := exec.execRunPreStart(t.Context(), Operation{Service: &service}) + assert.NilError(t, err) +} + +// TestExecRunPostStart_RunsHookAndStreamsToListener verifies that +// execRunPostStart resolves the replica the start chain just brought up and +// runs its post_start hooks, forwarding their output to the listener the +// executor was constructed with — the plumbing an attached `up` will use +// once the interactive-up convergence wires a real one in. +func TestExecRunPostStart_RunsHookAndStreamsToListener(t *testing.T) { + svc, apiClient, _ := newStartPhaseTestService(t) + + var mu sync.Mutex + var lines []string + listener := func(event api.ContainerEvent) { + mu.Lock() + defer mu.Unlock() + lines = append(lines, event.Line) + } + + serverConn, clientConn := net.Pipe() + go func() { + assert.NilError(t, writeStdcopyFrame(serverConn, 1, "post-start ok\n")) + serverConn.Close() //nolint:errcheck + }() + + apiClient.EXPECT().ExecCreate(gomock.Any(), "c1", gomock.Any()). + Return(client.ExecCreateResult{ID: "exec1"}, nil) + apiClient.EXPECT().ExecAttach(gomock.Any(), "exec1", gomock.Any()). + Return(client.ExecAttachResult{HijackedResponse: client.NewHijackedResponse(clientConn, "")}, nil) + apiClient.EXPECT().ExecInspect(gomock.Any(), "exec1", gomock.Any()). + Return(client.ExecInspectResult{ExitCode: 0}, nil) + + service := types.ServiceConfig{ + Name: "web", + PostStart: []types.ServiceHook{{Command: types.ShellCommand{"/notify.sh"}}}, + } + ctr := container.Summary{ID: "c1", Names: []string{"/test-web-1"}} + + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), listener) + err := exec.execRunPostStart(t.Context(), Operation{Service: &service, Container: &ctr}) + assert.NilError(t, err) + assert.DeepEqual(t, lines, []string{"post-start ok"}) +} + +// TestExecutePlanStartPhaseEventContract drives a full wait + enriched-start +// plan through the real executor and asserts the per-resource progressions +// the event contract promises (see api.EventProcessor): a dependency wait +// reports Waiting→Healthy on the dependency's own container, and the +// dependent's start reports a single Starting→Started on its own container — +// pre_start/post_start, when declared, fold into that same progression +// rather than surfacing their own. +func TestExecutePlanStartPhaseEventContract(t *testing.T) { + svc, apiClient, recorder := newStartPhaseTestService(t) + + project := &types.Project{ + Name: "test", + Services: types.Services{ + "db": {Name: "db"}, + "app": {Name: "app"}, + }, + } + dbSummary := container.Summary{ + ID: "db-id", + Names: []string{"/test-db-1"}, + Labels: map[string]string{api.ServiceLabel: "db", api.OneoffLabel: "False"}, + } + appSummary := container.Summary{ + ID: "app-id", + Names: []string{"/test-app-1"}, + Labels: map[string]string{api.ServiceLabel: "app", api.OneoffLabel: "False"}, + } + observed := emptyObservedState("test") + + apiClient.EXPECT().ContainerInspect(gomock.Any(), "db-id", gomock.Any()).Return(client.ContainerInspectResult{ + Container: container.InspectResponse{ + ID: "db-id", + Name: "/test-db-1", + State: &container.State{ + Status: container.StateRunning, + Health: &container.Health{Status: container.Healthy}, + }, + Config: &container.Config{Healthcheck: &container.HealthConfig{Test: []string{"CMD", "true"}}}, + }, + }, nil) + apiClient.EXPECT().ContainerStart(gomock.Any(), "app-id", gomock.Any()).Return(client.ContainerStartResult{}, nil) + + appSvc := project.Services["app"] + plan := &Plan{} + wait := plan.addNode(Operation{ + Type: OpWaitCondition, + ResourceID: "wait:db:service_healthy", + Cause: "app depends on db (service_healthy)", + Name: "db", + Condition: types.ServiceConditionHealthy, + }, "") + plan.addNode(Operation{ + Type: OpStartContainer, + ResourceID: "service:app:1", + Cause: "start", + Service: &appSvc, + Container: &appSummary, + }, "start:app:1", wait) + + exec := svc.newPlanExecutor(project, observed, nil) + exec.containersByService["db"] = Containers{dbSummary} + + assert.NilError(t, exec.run(t.Context(), plan)) + + assert.DeepEqual(t, recorder.byID["Container test-db-1"], []string{api.StatusWaiting, api.StatusHealthy}) + assert.DeepEqual(t, recorder.byID["Container test-app-1"], []string{api.StatusStarting, api.StatusStarted}) +} + +// TestExecutePlanStartPhaseFullChainThroughCreateNode drives +// RunPreStart→StartContainer→RunPostStart together through the real run() +// DAG/group-tracker path for a replica the plan itself creates (no observed +// container.Summary — CreateNodeID only), exercising resolveContainerSummary's +// live-view lookup. It asserts the group emits exactly one Starting→Started +// pair, with Starting firing at the StartContainer node — not at +// RunPreStart, which must run silently first, matching the imperative +// engine's startService. +func TestExecutePlanStartPhaseFullChainThroughCreateNode(t *testing.T) { + svc, apiClient, recorder := newStartPhaseTestService(t) + + service := types.ServiceConfig{ + Name: "web", + PreStart: []types.PreStartHook{{ContainerSpec: types.ContainerSpec{Command: types.ShellCommand{"init"}}}}, + PostStart: []types.ServiceHook{{Command: types.ShellCommand{"/notify.sh"}}}, + } + project := &types.Project{Name: "test", Services: types.Services{"web": service}} + + apiClient.EXPECT().ContainerCreate(gomock.Any(), gomock.Any()). + Return(client.ContainerCreateResult{ID: "new-id"}, nil) + apiClient.EXPECT().ContainerInspect(gomock.Any(), "new-id", gomock.Any()). + Return(client.ContainerInspectResult{Container: container.InspectResponse{ + ID: "new-id", + Name: "/test-web-1", + Config: &container.Config{}, + NetworkSettings: &container.NetworkSettings{}, + }}, nil) + // pre_start's daemon drift-check: no replica running yet, then the + // runner scan finds the hook-0 runner the (untested here) create phase + // prepared, and its wait→logs→start→remove sequence runs to success. + driftCheck := apiClient.EXPECT().ContainerList(gomock.Any(), gomock.Any()). + Return(client.ContainerListResult{}, nil) + scan := expectRunnerScan(apiClient, runnerSummary("hook-1", 0)).After(driftCheck) + wait := apiClient.EXPECT().ContainerWait(gomock.Any(), "hook-1", gomock.Any()). + Return(waitResultExit(0)).After(scan) + logs := apiClient.EXPECT().ContainerLogs(gomock.Any(), "hook-1", gomock.Any()). + Return(emptyLogs(), nil).After(wait) + hookStart := apiClient.EXPECT().ContainerStart(gomock.Any(), "hook-1", gomock.Any()). + Return(client.ContainerStartResult{}, nil).After(logs) + expectSuccessRemove(apiClient, "hook-1").After(hookStart) + + apiClient.EXPECT().ContainerStart(gomock.Any(), "new-id", gomock.Any()). + Return(client.ContainerStartResult{}, nil) + apiClient.EXPECT().ExecCreate(gomock.Any(), "new-id", gomock.Any()). + Return(client.ExecCreateResult{ID: "exec1"}, nil) + apiClient.EXPECT().ExecAttach(gomock.Any(), "exec1", gomock.Any()). + DoAndReturn(func(context.Context, string, client.ExecAttachOptions) (client.ExecAttachResult, error) { + serverConn, clientConn := net.Pipe() + go serverConn.Close() //nolint:errcheck + return client.ExecAttachResult{HijackedResponse: client.NewHijackedResponse(clientConn, "")}, nil + }) + apiClient.EXPECT().ExecInspect(gomock.Any(), "exec1", gomock.Any()). + Return(client.ExecInspectResult{ExitCode: 0}, nil) + + plan := &Plan{} + create := plan.addNode(Operation{ + Type: OpCreateContainer, + ResourceID: "service:web:1", + Cause: "no existing container", + Service: &service, + Name: "test-web-1", + Number: 1, + }, "") + preStart := plan.addNode(Operation{ + Type: OpRunPreStart, + ResourceID: "service:web:1", + Cause: "pre_start hooks", + Service: &service, + CreateNodeID: create.ID, + }, "start:web:1", create) + start := plan.addNode(Operation{ + Type: OpStartContainer, + ResourceID: "service:web:1", + Cause: "start", + Service: &service, + CreateNodeID: create.ID, + }, "start:web:1", preStart) + plan.addNode(Operation{ + Type: OpRunPostStart, + ResourceID: "service:web:1", + Cause: "post_start hooks", + Service: &service, + CreateNodeID: create.ID, + }, "start:web:1", start) + + exec := svc.newPlanExecutor(project, emptyObservedState("test"), nil) + assert.NilError(t, exec.run(t.Context(), plan)) + + // The create node is ungrouped (Creating/Created), then the start group + // contributes exactly one Starting/Started pair — pre_start and + // post_start run silently within it. + assert.DeepEqual(t, recorder.byID["Container test-web-1"], []string{"Creating", "Created", api.StatusStarting, api.StatusStarted}) +} + +// TestExecutePlanMissingRequiredDependencyFailsSilently verifies, through the +// real run()/emitErrorEvent path, that a required dependency with no live +// container at all fails the plan without emitting any spurious event — +// matching waitDependencies, which reports this case as a returned error +// only (see execWaitCondition's len(waitingFor)==0 branch). +func TestExecutePlanMissingRequiredDependencyFailsSilently(t *testing.T) { + svc, _, recorder := newStartPhaseTestService(t) + + plan := &Plan{} + plan.addNode(Operation{ + Type: OpWaitCondition, + ResourceID: "wait:db:service_healthy", + Cause: "app depends on db (service_healthy)", + Name: "db", + Condition: types.ServiceConditionHealthy, + }, "") + + err := svc.executePlan(t.Context(), &types.Project{Name: "test"}, emptyObservedState("test"), plan) + assert.ErrorContains(t, err, "missing dependency db") + assert.Equal(t, len(recorder.byID), 0, "a missing-dependency failure must not emit a spurious event") +} diff --git a/pkg/compose/executor_test.go b/pkg/compose/executor_test.go index e8da0c64f5c..3719e412ee6 100644 --- a/pkg/compose/executor_test.go +++ b/pkg/compose/executor_test.go @@ -297,7 +297,7 @@ func TestExecutePlanRemoveContainerDropsFromCache(t *testing.T) { Container: &oldCtr, }, "", stopNode) - exec := svc.newPlanExecutor(&types.Project{Name: "test"}, observed) + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, observed, nil) assert.NilError(t, exec.run(t.Context(), plan)) assert.Equal(t, len(exec.containersByService["web"]), 0, @@ -363,7 +363,7 @@ func TestExecutePlanConcurrentRemovesCacheCoherence(t *testing.T) { }, "", stop) } - exec := svc.newPlanExecutor(&types.Project{Name: "test"}, observed) + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, observed, nil) assert.NilError(t, exec.run(t.Context(), plan)) assert.Equal(t, len(exec.containersByService["web"]), 0, @@ -412,7 +412,7 @@ func TestExecutePlanRespectsMaxConcurrencyAcrossDependencyChain(t *testing.T) { }, "", deps...) } - exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test")) + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil) done := make(chan error, 1) go func() { done <- exec.run(t.Context(), plan) }() @@ -481,7 +481,7 @@ func TestExecutePlanIndependentNodeNotSerializedBehindADependencyWait(t *testing Type: OpStopContainer, ResourceID: "service:app:1", Cause: "unrelated", Container: &independent, }, "") - exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test")) + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil) done := make(chan error, 1) go func() { done <- exec.run(t.Context(), plan) }() @@ -674,7 +674,7 @@ func TestExecutePlanCreateNetworkConflictIsSuccess(t *testing.T) { func TestExecRemoveNetworkBestEffort(t *testing.T) { newExec := func(t *testing.T) (*planExecutor, *mocks.MockAPIClient) { svc, apiClient := newTestService(t) - return svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test")), apiClient + return svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil), apiClient } t.Run("best-effort ignores conflict", func(t *testing.T) { From 0c39d426bdb2a554999d04bf769cc02aca4f8e16 Mon Sep 17 00:00:00 2001 From: Nicolas De Loof Date: Wed, 30 Sep 2026 14:53:30 +0200 Subject: [PATCH 2/5] fix: a failed pre_start/start no longer lets its dependent run A node's done-channel closes strictly before errgroup cancels its derived context on that same error - the close happens synchronously inside the failing node's goroutine, ctx cancellation only after it returns. A dependent blocked in run()'s select { case <-done[dep.ID]: ...; case <-ctx.Done(): ... } therefore always sees its dependency's done-channel ready first, regardless of whether that dependency failed for a genuine reason or was cancelled. execCreateHookContainer already guarded against exactly this. The two other new dispatch functions from the previous commit did not: execStartContainer's enriched branch could call ContainerStart right after its OpRunPreStart failed for real (e.g. "no hook runner container found"), contradicting pre_start.go's documented guarantee that a non-zero hook exit gates service start - and symmetrically for execRunPostStart after a failed OpStartContainer. Both now bail out on ctx.Err() first, mirroring the existing guard. Pinned by a regression test reproducing the race deterministically: no ContainerStart expectation is registered, so the unguarded code calls it from inside the errgroup worker goroutine, where gomock's missing- expectation failure hangs the test binary instead of failing it cleanly - a stronger tell than a plain assertion failure would have been. Signed-off-by: Nicolas De Loof --- pkg/compose/executor_ops.go | 16 +++++++ pkg/compose/executor_start_test.go | 71 ++++++++++++++++++++++++++++++ 2 files changed, 87 insertions(+) diff --git a/pkg/compose/executor_ops.go b/pkg/compose/executor_ops.go index 9b303109110..79433775670 100644 --- a/pkg/compose/executor_ops.go +++ b/pkg/compose/executor_ops.go @@ -144,6 +144,14 @@ func (exec *planExecutor) execStartContainer(ctx context.Context, op Operation) return err } + // A dependency's done-channel also closes on failure: when a preceding + // node of this replica's chain (a wait, pre_start) failed and canceled + // the run, bail out on the cancellation instead of starting a container + // whose pre_start hook just failed — same race as execCreateHookContainer. + if err := ctx.Err(); err != nil { + return err + } + id, err := exec.resolveContainerID(op) if err != nil { return err @@ -262,6 +270,14 @@ func (exec *planExecutor) execRunPreStart(ctx context.Context, op Operation) err // execRunPostStart runs the service's post_start hooks against the replica // the start chain just brought up. func (exec *planExecutor) execRunPostStart(ctx context.Context, op Operation) error { + // Same race as execCreateHookContainer/execStartContainer: a failed + // StartContainer's done-channel closes before errgroup's ctx is + // canceled, so this must not run post_start against a container that + // never actually started. + if err := ctx.Err(); err != nil { + return err + } + ctr, err := exec.resolveContainerSummary(op) if err != nil { return err diff --git a/pkg/compose/executor_start_test.go b/pkg/compose/executor_start_test.go index 2e61707c241..27852ddad22 100644 --- a/pkg/compose/executor_start_test.go +++ b/pkg/compose/executor_start_test.go @@ -435,3 +435,74 @@ func TestExecutePlanMissingRequiredDependencyFailsSilently(t *testing.T) { assert.ErrorContains(t, err, "missing dependency db") assert.Equal(t, len(recorder.byID), 0, "a missing-dependency failure must not emit a spurious event") } + +// TestExecutePlanFailedPreStartGatesStart is a regression test for a race in +// run()'s DAG walk: a node's done-channel closes (unblocking its dependents) +// strictly before errgroup cancels its derived context on that same error — +// the close happens synchronously inside the failing node's goroutine, ctx +// cancellation only after it returns. A dependent blocked in run()'s +// `select { case <-done[dep.ID]: ...; case <-ctx.Done(): ... }` therefore +// always sees its dependency's done-channel ready first, deterministically, +// regardless of whether that dependency failed for a genuine reason or was +// itself cancelled. +// +// execCreateHookContainer already guards against this with a leading +// ctx.Err() check; this test pins the same guard on execStartContainer's +// enriched branch: when OpRunPreStart fails for a real reason (here: no hook +// runner container found — never a context cancellation), the OpStartContainer +// depending on it must not call ContainerStart. No ContainerStart expectation +// is registered below: an unexpected call fails the test. +func TestExecutePlanFailedPreStartGatesStart(t *testing.T) { + svc, apiClient, _ := newStartPhaseTestService(t) + + service := types.ServiceConfig{ + Name: "web", + PreStart: []types.PreStartHook{{ContainerSpec: types.ContainerSpec{Command: types.ShellCommand{"init"}}}}, + } + project := &types.Project{Name: "test", Services: types.Services{"web": service}} + + apiClient.EXPECT().ContainerCreate(gomock.Any(), gomock.Any()). + Return(client.ContainerCreateResult{ID: "new-id"}, nil) + apiClient.EXPECT().ContainerInspect(gomock.Any(), "new-id", gomock.Any()). + Return(client.ContainerInspectResult{Container: container.InspectResponse{ + ID: "new-id", + Name: "/test-web-1", + Config: &container.Config{}, + NetworkSettings: &container.NetworkSettings{}, + }}, nil) + // Both the pre_start drift-check and the hook-runner scan list no + // container: no replica is running, and no runner was prepared for the + // declared hook — runPreStart fails with "no hook runner container + // found", a genuine error, not a cancellation. + apiClient.EXPECT().ContainerList(gomock.Any(), gomock.Any()). + Return(client.ContainerListResult{}, nil).Times(2) + // No ContainerStart expectation: registering one here would hide the + // bug by silently accepting the call this test exists to forbid. + + plan := &Plan{} + create := plan.addNode(Operation{ + Type: OpCreateContainer, + ResourceID: "service:web:1", + Cause: "no existing container", + Service: &service, + Name: "test-web-1", + Number: 1, + }, "") + preStart := plan.addNode(Operation{ + Type: OpRunPreStart, + ResourceID: "service:web:1", + Cause: "pre_start hooks", + Service: &service, + CreateNodeID: create.ID, + }, "start:web:1", create) + plan.addNode(Operation{ + Type: OpStartContainer, + ResourceID: "service:web:1", + Cause: "start", + Service: &service, + CreateNodeID: create.ID, + }, "start:web:1", preStart) + + err := svc.executePlan(t.Context(), project, emptyObservedState("test"), plan) + assert.ErrorContains(t, err, "no hook runner container found") +} From 219f34cba50c86705daae19df534e22c46ef5ee1 Mon Sep 17 00:00:00 2001 From: Nicolas De Loof Date: Wed, 30 Sep 2026 15:33:44 +0200 Subject: [PATCH 3/5] docs: spell out that the ctx.Err() guards narrow the race, not close it docker-agent's review of the previous commit's fix correctly pointed out that the ctx.Err() guard added to execStartContainer/execRunPostStart narrows the race it targets but does not eliminate it: close(done[...]) still runs inside the failing node's own goroutine, strictly before errgroup's cancel() on that goroutine's returned error. A dependent unblocked from <-done[dep.ID] can observe ctx.Err() == nil for a brief window even though its dependency just failed for a genuine reason. This is exactly the structural gap the epic (#14081) already flagged in its "failed-dependency semantics" comment before this PR existed: the real fix carries each node's success/failure through what dependents wait on (e.g. a chan error) instead of inferring it from ctx.Err(), uniformly across every operation type - including the already-merged execCreateHookContainer, which has the same limit. That is a dedicated brick of the executor lot, not squeezed into this PR's diff. No behavior change: this documents the residual window explicitly at both guard sites and in the regression test, so the next reader (and the next PR) doesn't mistake "narrows" for "closes". Signed-off-by: Nicolas De Loof --- pkg/compose/executor_ops.go | 26 ++++++++++++++++++-------- pkg/compose/executor_start_test.go | 18 +++++++++--------- 2 files changed, 27 insertions(+), 17 deletions(-) diff --git a/pkg/compose/executor_ops.go b/pkg/compose/executor_ops.go index 79433775670..47192c2d9c2 100644 --- a/pkg/compose/executor_ops.go +++ b/pkg/compose/executor_ops.go @@ -144,10 +144,19 @@ func (exec *planExecutor) execStartContainer(ctx context.Context, op Operation) return err } - // A dependency's done-channel also closes on failure: when a preceding - // node of this replica's chain (a wait, pre_start) failed and canceled - // the run, bail out on the cancellation instead of starting a container - // whose pre_start hook just failed — same race as execCreateHookContainer. + // A dependency's done-channel closes (unblocking this node) before + // errgroup cancels ctx on that same dependency's error — cancel() only + // runs after the failing goroutine returns, close(done[...]) runs inside + // it. This check narrows that window (same guard as + // execCreateHookContainer) but cannot close it: a preceding node of this + // replica's chain (a wait, pre_start) can fail for a genuine reason a + // moment before this goroutine observes ctx.Err(), still nil, and starts + // a container whose pre_start hook just failed. The structural fix — + // carrying each node's success/failure through what dependents wait on, + // instead of inferring it from ctx.Err() — is tracked as its own brick of + // the executor lot (#14081, see the epic's "failed-dependency semantics" + // comment); every per-operation ctx.Err() guard here, including this one, + // is an interim narrowing, not the fix. if err := ctx.Err(); err != nil { return err } @@ -270,10 +279,11 @@ func (exec *planExecutor) execRunPreStart(ctx context.Context, op Operation) err // execRunPostStart runs the service's post_start hooks against the replica // the start chain just brought up. func (exec *planExecutor) execRunPostStart(ctx context.Context, op Operation) error { - // Same race as execCreateHookContainer/execStartContainer: a failed - // StartContainer's done-channel closes before errgroup's ctx is - // canceled, so this must not run post_start against a container that - // never actually started. + // Same interim narrowing as execStartContainer's guard above, not a full + // fix (see its comment): a failed StartContainer's done-channel closes + // before errgroup's ctx is canceled, so a narrow window remains where + // this could still run post_start against a container whose start just + // failed for a genuine reason. if err := ctx.Err(); err != nil { return err } diff --git a/pkg/compose/executor_start_test.go b/pkg/compose/executor_start_test.go index 27852ddad22..2bd37ae3b56 100644 --- a/pkg/compose/executor_start_test.go +++ b/pkg/compose/executor_start_test.go @@ -438,16 +438,16 @@ func TestExecutePlanMissingRequiredDependencyFailsSilently(t *testing.T) { // TestExecutePlanFailedPreStartGatesStart is a regression test for a race in // run()'s DAG walk: a node's done-channel closes (unblocking its dependents) -// strictly before errgroup cancels its derived context on that same error — -// the close happens synchronously inside the failing node's goroutine, ctx -// cancellation only after it returns. A dependent blocked in run()'s -// `select { case <-done[dep.ID]: ...; case <-ctx.Done(): ... }` therefore -// always sees its dependency's done-channel ready first, deterministically, -// regardless of whether that dependency failed for a genuine reason or was -// itself cancelled. +// inside the failing node's own goroutine, strictly before errgroup calls +// cancel() on that goroutine's returned error — cancel() only runs after it +// returns. A dependent unblocked from `<-done[dep.ID]` can therefore observe +// ctx.Err() == nil for a brief window even though its dependency just failed +// for a genuine reason, not a cancellation. // -// execCreateHookContainer already guards against this with a leading -// ctx.Err() check; this test pins the same guard on execStartContainer's +// execCreateHookContainer already narrows this with a leading ctx.Err() +// check (it cannot close the window — see the guard's comment in +// execStartContainer for why this is deferred to a dedicated executor-lot +// fix, #14081); this test pins the same narrowing on execStartContainer's // enriched branch: when OpRunPreStart fails for a real reason (here: no hook // runner container found — never a context cancellation), the OpStartContainer // depending on it must not call ContainerStart. No ContainerStart expectation From b1a06a5de263032427d75b8c98595d62eafcf864 Mon Sep 17 00:00:00 2001 From: Nicolas De Loof Date: Thu, 1 Oct 2026 16:02:33 +0200 Subject: [PATCH 4/5] fix: execRunPreStart bails out on a canceled context Same interim narrowing as execStartContainer and execRunPostStart: the node hangs off the replica's create node, so when that create fails the pre_start hooks could still fire in the window before errgroup cancels the context. The structural fix stays tracked in #14081. Signed-off-by: Nicolas De Loof --- pkg/compose/executor_ops.go | 9 +++++++++ pkg/compose/executor_start_test.go | 20 ++++++++++++++++++++ 2 files changed, 29 insertions(+) diff --git a/pkg/compose/executor_ops.go b/pkg/compose/executor_ops.go index 47192c2d9c2..2231cd350b1 100644 --- a/pkg/compose/executor_ops.go +++ b/pkg/compose/executor_ops.go @@ -265,6 +265,15 @@ func (exec *planExecutor) execWaitCondition(ctx context.Context, op Operation) e // accepts — skips a second, redundant run of the hooks rather than erroring // on runners already consumed. func (exec *planExecutor) execRunPreStart(ctx context.Context, op Operation) error { + // Same interim narrowing as execStartContainer's guard (see its + // comment), not a full fix: this node hangs off the replica's create + // node, whose done-channel closes before errgroup's ctx is canceled when + // that create fails for a genuine reason, so a narrow window remains + // where pre_start hooks could still run against the runner containers. + if err := ctx.Err(); err != nil { + return err + } + running, err := exec.compose.getContainers(ctx, exec.project.Name, oneOffExclude, false, op.Service.Name) if err != nil { return err diff --git a/pkg/compose/executor_start_test.go b/pkg/compose/executor_start_test.go index 2bd37ae3b56..3f89f59d747 100644 --- a/pkg/compose/executor_start_test.go +++ b/pkg/compose/executor_start_test.go @@ -211,6 +211,26 @@ func TestExecRunPreStart_SkipsWhenReplicaAlreadyRunning(t *testing.T) { assert.NilError(t, err) } +// TestExecRunPreStart_BailsOutWhenContextCanceled pins the same interim +// narrowing execStartContainer and execRunPostStart carry: once ctx is +// canceled (a sibling node failed), the node must neither touch the daemon +// nor run hooks. No ContainerList expectation is registered: an unexpected +// call fails the test. +func TestExecRunPreStart_BailsOutWhenContextCanceled(t *testing.T) { + svc, _, _ := newStartPhaseTestService(t) + + service := types.ServiceConfig{ + Name: "web", + PreStart: []types.PreStartHook{{ContainerSpec: types.ContainerSpec{Command: types.ShellCommand{"init"}}}}, + } + ctx, cancel := context.WithCancel(t.Context()) + cancel() + + exec := svc.newPlanExecutor(&types.Project{Name: "test"}, emptyObservedState("test"), nil) + err := exec.execRunPreStart(ctx, Operation{Service: &service}) + assert.ErrorIs(t, err, context.Canceled) +} + // TestExecRunPostStart_RunsHookAndStreamsToListener verifies that // execRunPostStart resolves the replica the start chain just brought up and // runs its post_start hooks, forwarding their output to the listener the From de7ce858ad1358d6828417febf164c5b3abf3add Mon Sep 17 00:00:00 2001 From: Nicolas De Loof Date: Fri, 2 Oct 2026 09:01:40 +0200 Subject: [PATCH 5/5] fix: release startMx with defer in the enriched execStartContainer path Signed-off-by: Nicolas De Loof --- pkg/compose/executor_ops.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/compose/executor_ops.go b/pkg/compose/executor_ops.go index 2231cd350b1..6f5fb6793d3 100644 --- a/pkg/compose/executor_ops.go +++ b/pkg/compose/executor_ops.go @@ -173,8 +173,8 @@ func (exec *planExecutor) execStartContainer(ctx context.Context, op Operation) } startMx.Lock() + defer startMx.Unlock() _, err = exec.compose.apiClient().ContainerStart(ctx, id, client.ContainerStartOptions{}) - startMx.Unlock() return err }