Skip to content
Open
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
22 changes: 22 additions & 0 deletions pkg/api/event.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
26 changes: 24 additions & 2 deletions pkg/compose/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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(),
}
Expand Down Expand Up @@ -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:
Expand Down
114 changes: 95 additions & 19 deletions pkg/compose/executor_events.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[medium] A failing OpRunPreStart emits an api.Error event with no preceding Working/Starting event

triggersWorking returns true only for OpStartContainer in a groupStart group:

func (gs *groupState) triggersWorking(t OperationType) bool {
    if gs.kind == groupStart {
        return t == OpStartContainer  // OpRunPreStart returns false
    }
    return gs.started == 1
}

This is intentional for the happy path — pre_start hooks run silently before the Starting event fires. But when OpRunPreStart itself fails (e.g. the hook-runner container is not found, or the hook exits non-zero), run() calls onNodeError, which emits an api.Error resource event on gs.eventName. No Working event (status Starting) was ever emitted for this group, because triggersWorking returned false for OpRunPreStart.

Progress-bar consumers tracking Working → Done / Error state transitions will see a bare Error for a container they had never seen transition to Starting. The resource appears in an error state from nowhere, which can cause progress-display UIs to render incorrectly (the container never appeared "starting" before flipping to an error).

The imperative engine's startService always emits a Starting event before any outcome; the groupStart / onNodeError path has no equivalent guard.

Suggested fix: track whether the group's Working event has fired (e.g. a workingEmitted bool field on groupState), and emit an implicit Starting event in onNodeError when the group kind is groupStart and !gs.workingEmitted — keeping the event sequence well-formed even when pre_start fails before ContainerStart is reached.

Confidence Score
🟡 moderate 75/100

}
return gs.started == 1
}

func (gt *groupTracker) onNodeDone(node *PlanNode, events api.EventProcessor) {
Expand All @@ -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()))
}
}

Expand Down Expand Up @@ -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:
Expand Down
Loading
Loading