diff --git a/internal/coordinator/coordinator.go b/internal/coordinator/coordinator.go new file mode 100644 index 0000000000..fb1fef379e --- /dev/null +++ b/internal/coordinator/coordinator.go @@ -0,0 +1,182 @@ +/* + 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 coordinator holds the coordinator-specific integration used by +// Compose: detecting a coordinator-enabled Docker context and pushing the +// project configuration to it over the engine socket. It mirrors the +// self-contained layout of internal/desktop (detection, HTTP client, versioned +// URL and outcome handling) so that pkg/compose keeps only a thin call-site. +package coordinator + +import ( + "bytes" + "context" + "fmt" + "io" + "net" + "net/http" + "strconv" + "strings" + "time" + + "github.com/compose-spec/compose-go/v2/types" + "github.com/docker/cli/cli/command" + "github.com/moby/moby/client/pkg/versions" +) + +// MetadataKey is the Docker context metadata key that identifies a coordinator +// behind a Docker API socket. When the current context's metadata carries this +// key set to a truthy value, the socket is assumed to accept POST +// /v{version}/compose/project (see MinAPIVersion), and Compose sends the +// project configuration to the coordinator at the start of "compose up" before +// any other Docker API call. +const MetadataKey = "com.docker.compose.coordinator" + +// MinAPIVersion is the Engine/coordinator API version that introduced +// POST /v{version}/compose/project for the Compose project-config push. It +// documents the minimum coordinator version and is the single place the push +// gates on. +const MinAPIVersion = "1.51" + +// defaultTimeout bounds a single project-config push. The push runs before any +// other Docker API call in "compose up", so without a deadline a coordinator +// that accepts the connection but never responds would hang "up" indefinitely, +// defeating the non-fatal "warn and continue" guarantee. +const defaultTimeout = 10 * time.Second + +// responseBodyLimit bounds how much of an error response body is read back +// into the returned error message. +const responseBodyLimit = 2048 + +// CompleteHeader signals whether the pushed project is the whole project or a +// subset. "compose up" with no service arguments resolves the entire project +// and sends "true"; "compose up " narrows the project to the named +// services plus their dependency closure and sends "false", telling the +// coordinator to merge the payload with previously pushed config rather than +// treat it as authoritative. An absent header (older Compose clients) must be +// read as "false": those clients may also have pushed a subset, so the +// coordinator must never prune on their behalf. +const CompleteHeader = "X-Compose-Project-Complete" + +// Enabled reports whether the current Docker context represents a compose +// coordinator via the MetadataKey metadata key. A metadata read failure is +// treated as "not enabled" so that "compose up" is never blocked on context +// inspection. The value is accepted as either a JSON boolean true or the string +// "true" (case-insensitive), mirroring how custom context metadata may be +// decoded (see ConfigFromDockerContext in internal/tracing/docker_context.go). +func Enabled(dockerCli command.Cli) bool { + meta, err := dockerCli.ContextStore().GetMetadata(dockerCli.CurrentContext()) + if err != nil { + return false + } + + var value any + switch m := meta.Metadata.(type) { + case command.DockerContext: + value = m.AdditionalFields[MetadataKey] + case map[string]any: + value = m[MetadataKey] + } + + switch v := value.(type) { + case bool: + return v + case string: + return strings.EqualFold(v, "true") + default: + return false + } +} + +// Client pushes the Compose project configuration to a coordinator over the +// Docker engine socket, using the engine dialer for transport (mirroring +// internal/desktop/client.go). +type Client struct { + dialer func(ctx context.Context) (net.Conn, error) + timeout time.Duration +} + +// NewClient builds a coordinator client that reaches the engine over the given +// dialer (typically apiClient.Dialer()). +func NewClient(dialer func(ctx context.Context) (net.Conn, error)) *Client { + return &Client{dialer: dialer, timeout: defaultTimeout} +} + +// PushProjectConfig sends the full project configuration to the coordinator's +// version-negotiated POST /v{version}/compose/project endpoint. apiVersion is +// the negotiated engine API version; the push is rejected before any network +// call when it predates MinAPIVersion. A non-2xx response is returned as an +// error; callers warn and continue. The request is bounded by a timeout so a +// coordinator that never responds cannot hang "compose up". +// +// complete reports whether project is the whole project (true) or a subset +// that the coordinator should merge with what it already holds (false); it is +// conveyed via CompleteHeader. +func (c *Client) PushProjectConfig(ctx context.Context, apiVersion string, project *types.Project, complete bool) error { + if versions.LessThan(apiVersion, MinAPIVersion) { + return fmt.Errorf("coordinator API version %s does not support the project-config push (requires %s or later)", + apiVersion, MinAPIVersion) + } + + payload, err := project.MarshalJSON() + if err != nil { + return fmt.Errorf("marshaling project config: %w", err) + } + + httpClient := &http.Client{ + Timeout: c.timeout, + Transport: &http.Transport{ + DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { + return c.dialer(ctx) + }, + }, + } + + // The host is cosmetic: the custom dialer handles routing to the engine + // socket. It exists only to form a valid URL (see desktop.backendURL). + url := fmt.Sprintf("http://docker/v%s/compose/project", apiVersion) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(payload)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + // Always set the header explicitly (not omit-on-false) so the value is + // unambiguous on the wire; only absence means "older client". + req.Header.Set(CompleteHeader, strconv.FormatBool(complete)) + + resp, err := httpClient.Do(req) + if err != nil { + return err + } + defer func() { + _ = resp.Body.Close() + // The transport is discarded after this call; tear down its idle + // persistConn goroutines eagerly instead of waiting out + // IdleConnTimeout, which matters when "up" is invoked in a loop. + httpClient.CloseIdleConnections() + }() + + if resp.StatusCode >= http.StatusBadRequest { + body, _ := io.ReadAll(io.LimitReader(resp.Body, responseBodyLimit)) + msg := strings.TrimSpace(string(body)) + if msg == "" { + return fmt.Errorf("coordinator returned status %d", resp.StatusCode) + } + return fmt.Errorf("coordinator returned status %d: %s", resp.StatusCode, msg) + } + + return nil +} diff --git a/internal/coordinator/coordinator_test.go b/internal/coordinator/coordinator_test.go new file mode 100644 index 0000000000..9f0dd57b76 --- /dev/null +++ b/internal/coordinator/coordinator_test.go @@ -0,0 +1,248 @@ +/* + 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 coordinator + +import ( + "context" + "errors" + "io" + "net" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/compose-spec/compose-go/v2/types" + "github.com/docker/cli/cli/command" + "github.com/docker/cli/cli/context/store" + "go.uber.org/mock/gomock" + "gotest.tools/v3/assert" + + "github.com/docker/compose/v5/pkg/mocks" +) + +var testStoreCfg = store.NewConfig( + func() any { + return &map[string]any{} + }, +) + +func newContextStore(t *testing.T, meta any) store.Store { + t.Helper() + st := store.New(t.TempDir(), testStoreCfg) + err := st.CreateOrUpdate(store.Metadata{ + Name: "test", + Metadata: meta, + Endpoints: make(map[string]any), + }) + assert.NilError(t, err) + return st +} + +func TestPushEnabled(t *testing.T) { + if testing.Short() { + t.Skip("Requires filesystem access") + } + + dockerContext := func(fields map[string]any) command.DockerContext { + return command.DockerContext{Description: "test", AdditionalFields: fields} + } + + tests := []struct { + name string + meta any + want bool + }{ + { + name: "boolean true", + meta: dockerContext(map[string]any{MetadataKey: true}), + want: true, + }, + { + name: "string true", + meta: dockerContext(map[string]any{MetadataKey: "true"}), + want: true, + }, + { + name: "string TRUE case-insensitive", + meta: dockerContext(map[string]any{MetadataKey: "TRUE"}), + want: true, + }, + { + name: "boolean false", + meta: dockerContext(map[string]any{MetadataKey: false}), + want: false, + }, + { + name: "string other", + meta: dockerContext(map[string]any{MetadataKey: "yes"}), + want: false, + }, + { + name: "key absent", + meta: dockerContext(map[string]any{"other": true}), + want: false, + }, + { + name: "raw map form", + meta: map[string]any{MetadataKey: true}, + want: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + mockCtrl := gomock.NewController(t) + cli := mocks.NewMockCli(mockCtrl) + cli.EXPECT().ContextStore().Return(newContextStore(t, tt.meta)).AnyTimes() + cli.EXPECT().CurrentContext().Return("test").AnyTimes() + + assert.Equal(t, Enabled(cli), tt.want) + }) + } +} + +func TestPushEnabledMetadataError(t *testing.T) { + mockCtrl := gomock.NewController(t) + cli := mocks.NewMockCli(mockCtrl) + // An empty store returns an error for an unknown context; that must be + // treated as "not enabled" rather than blocking up. + cli.EXPECT().ContextStore().Return(store.New(t.TempDir(), testStoreCfg)).AnyTimes() + cli.EXPECT().CurrentContext().Return("missing").AnyTimes() + + assert.Equal(t, Enabled(cli), false) +} + +func newTestProject() *types.Project { + return &types.Project{ + Name: "test", + Services: types.Services{ + "web": {Name: "web", Image: "nginx"}, + }, + } +} + +func dialerFor(addr string) func(context.Context) (net.Conn, error) { + return func(ctx context.Context) (net.Conn, error) { + var d net.Dialer + return d.DialContext(ctx, "tcp", addr) + } +} + +func TestPushProjectConfig(t *testing.T) { + var ( + gotMethod string + gotPath string + gotBody []byte + gotComplete string + ) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotMethod = r.Method + gotPath = r.URL.Path + gotBody, _ = io.ReadAll(r.Body) + gotComplete = r.Header.Get(CompleteHeader) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + c := NewClient(dialerFor(srv.Listener.Addr().String())) + err := c.PushProjectConfig(t.Context(), MinAPIVersion, newTestProject(), true) + assert.NilError(t, err) + assert.Equal(t, gotMethod, http.MethodPost) + assert.Equal(t, gotPath, "/v"+MinAPIVersion+"/compose/project") + assert.Assert(t, len(gotBody) > 0) + assert.Equal(t, gotComplete, "true") +} + +func TestPushProjectConfigSubsetHeader(t *testing.T) { + // A subset push (empty-selection == false) must advertise itself so the + // coordinator merges rather than prunes. + var gotComplete string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotComplete = r.Header.Get(CompleteHeader) + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + c := NewClient(dialerFor(srv.Listener.Addr().String())) + err := c.PushProjectConfig(t.Context(), MinAPIVersion, newTestProject(), false) + assert.NilError(t, err) + assert.Equal(t, gotComplete, "false") +} + +func TestPushProjectConfigServerError(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "placement failed", http.StatusInternalServerError) + })) + defer srv.Close() + + c := NewClient(dialerFor(srv.Listener.Addr().String())) + err := c.PushProjectConfig(t.Context(), MinAPIVersion, newTestProject(), true) + assert.ErrorContains(t, err, "500") + assert.ErrorContains(t, err, "placement failed") +} + +func TestPushProjectConfigDialerError(t *testing.T) { + // A dialer that never connects surfaces as a request error, which callers + // treat as non-fatal. + dialer := func(context.Context) (net.Conn, error) { + return nil, errors.New("boom: no engine socket") + } + c := NewClient(dialer) + err := c.PushProjectConfig(t.Context(), MinAPIVersion, newTestProject(), true) + assert.ErrorContains(t, err, "boom: no engine socket") +} + +func TestPushProjectConfigErrorEmptyBody(t *testing.T) { + // A non-2xx status with no body still yields a useful error. + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusServiceUnavailable) + })) + defer srv.Close() + + c := NewClient(dialerFor(srv.Listener.Addr().String())) + err := c.PushProjectConfig(t.Context(), MinAPIVersion, newTestProject(), true) + assert.ErrorContains(t, err, "coordinator returned status 503") +} + +func TestPushProjectConfigVersionTooLow(t *testing.T) { + // The dialer must never be reached: the version gate rejects first. + dialer := func(context.Context) (net.Conn, error) { + t.Fatal("dialer should not be called when the API version is too low") + return nil, nil + } + c := NewClient(dialer) + err := c.PushProjectConfig(t.Context(), "1.44", newTestProject(), true) + assert.ErrorContains(t, err, "does not support the project-config push") +} + +func TestPushProjectConfigTimeout(t *testing.T) { + // A coordinator that accepts the connection but never responds must not + // hang the push: the client timeout bounds the request. + block := make(chan struct{}) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + <-block + })) + defer srv.Close() + defer close(block) + + c := NewClient(dialerFor(srv.Listener.Addr().String())) + c.timeout = 50 * time.Millisecond + + err := c.PushProjectConfig(t.Context(), MinAPIVersion, newTestProject(), true) + assert.ErrorContains(t, err, "context deadline exceeded") +} diff --git a/pkg/compose/up.go b/pkg/compose/up.go index b2bcb3f7db..530477862e 100644 --- a/pkg/compose/up.go +++ b/pkg/compose/up.go @@ -36,6 +36,7 @@ import ( "golang.org/x/sync/errgroup" "github.com/docker/compose/v5/cmd/formatter" + "github.com/docker/compose/v5/internal/coordinator" "github.com/docker/compose/v5/internal/desktop" "github.com/docker/compose/v5/internal/tracing" "github.com/docker/compose/v5/pkg/api" @@ -43,6 +44,22 @@ import ( func (s *composeService) Up(ctx context.Context, project *types.Project, options api.UpOptions) error { //nolint:gocyclo err := Run(ctx, tracing.SpanWrapFunc("project/up", tracing.ProjectOptions(ctx, project), func(ctx context.Context) error { + // When the current Docker context opts into the project-config push, send + // the project configuration to the coordinator before any other Docker API + // call. Failures are non-fatal: warn and continue bringing the project up. + if !s.dryRun && coordinator.Enabled(s.dockerCli) { + // The push is complete only when "up" resolved the whole project: + // an empty service selection AND no profiles applied. A non-empty + // selection narrows the project to a subset (+ dependency closure), + // and an applied profile disables the services outside it; either + // way the payload is partial and the coordinator must merge it + // rather than treat it as the authoritative project. + complete := len(options.Start.Services) == 0 && !hasActiveProfiles(project.Profiles) + if err := s.pushProjectConfig(ctx, project, complete); err != nil { + s.events.On(newEvent(api.ResourceCompose, api.Warning, "project config push to coordinator failed, continuing", err.Error())) + } + } + err := s.create(ctx, project, options.Create) if err != nil { return err @@ -302,6 +319,31 @@ func (s *composeService) Up(ctx context.Context, project *types.Project, options return err } +// pushProjectConfig negotiates the engine API version and hands the project +// off to the coordinator client, which owns the coordinator-specific transport +// and outcome handling (see internal/coordinator). complete reports whether +// project is the whole project or a subset the coordinator should merge. +// hasActiveProfiles reports whether any profile is applied to the project. +// project.Profiles carries a single empty-string entry when no profile is +// selected (see compose-go WithDefaultProfiles splitting an unset +// COMPOSE_PROFILES), so blank entries are ignored. +func hasActiveProfiles(profiles []string) bool { + for _, p := range profiles { + if p != "" { + return true + } + } + return false +} + +func (s *composeService) pushProjectConfig(ctx context.Context, project *types.Project, complete bool) error { + version, err := s.RuntimeAPIVersion(ctx) + if err != nil { + return fmt.Errorf("negotiating API version: %w", err) + } + return coordinator.NewClient(s.apiClient().Dialer()).PushProjectConfig(ctx, version, project, complete) +} + func shouldFollowStartEvent(event api.ContainerEvent, attached []string, attachTo []string) bool { if event.Type != api.ContainerEventStarted { return false diff --git a/pkg/compose/up_test.go b/pkg/compose/up_test.go index f38a9e1aee..2f6401a845 100644 --- a/pkg/compose/up_test.go +++ b/pkg/compose/up_test.go @@ -17,13 +17,115 @@ package compose import ( + "context" + "errors" + "net" + "net/http" + "net/http/httptest" "testing" + "github.com/compose-spec/compose-go/v2/types" + "github.com/moby/moby/client" + "go.uber.org/mock/gomock" "gotest.tools/v3/assert" + "github.com/docker/compose/v5/internal/coordinator" "github.com/docker/compose/v5/pkg/api" + "github.com/docker/compose/v5/pkg/mocks" ) +func newPushTestService(t *testing.T, apiClient *mocks.MockAPIClient, version string) *composeService { + t.Helper() + mockCtrl := gomock.NewController(t) + t.Cleanup(mockCtrl.Finish) + cli := mocks.NewMockCli(mockCtrl) + cli.EXPECT().Client().Return(apiClient).AnyTimes() + apiClient.EXPECT().Ping(gomock.Any(), client.PingOptions{NegotiateAPIVersion: true}). + Return(client.PingResult{APIVersion: version}, nil).AnyTimes() + apiClient.EXPECT().ClientVersion().Return(version).AnyTimes() + tested, err := NewComposeService(cli) + assert.NilError(t, err) + return tested.(*composeService) +} + +func dialerFor(addr string) func(context.Context) (net.Conn, error) { + return func(ctx context.Context) (net.Conn, error) { + var d net.Dialer + return d.DialContext(ctx, "tcp", addr) + } +} + +func newUpTestProject() *types.Project { + return &types.Project{ + Name: "test", + Services: types.Services{ + "web": {Name: "web", Image: "nginx"}, + }, + } +} + +// TestPushProjectConfigGlue covers the pkg/compose glue: it negotiates the +// engine API version and delegates to the coordinator client, which pushes to +// the version-negotiated endpoint reached over the engine dialer. +func TestPushProjectConfigGlue(t *testing.T) { + var gotPath string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotPath = r.URL.Path + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + mockCtrl := gomock.NewController(t) + apiClient := mocks.NewMockAPIClient(mockCtrl) + apiClient.EXPECT().Dialer().Return(dialerFor(srv.Listener.Addr().String())).AnyTimes() + s := newPushTestService(t, apiClient, coordinator.MinAPIVersion) + + err := s.pushProjectConfig(t.Context(), newUpTestProject(), true) + assert.NilError(t, err) + assert.Equal(t, gotPath, "/v"+coordinator.MinAPIVersion+"/compose/project") +} + +// TestPushProjectConfigGlueVersionError covers the glue's error branch when +// API-version negotiation fails; the coordinator dialer is never reached. +func TestPushProjectConfigGlueVersionError(t *testing.T) { + mockCtrl := gomock.NewController(t) + t.Cleanup(mockCtrl.Finish) + apiClient := mocks.NewMockAPIClient(mockCtrl) + cli := mocks.NewMockCli(mockCtrl) + cli.EXPECT().Client().Return(apiClient).AnyTimes() + apiClient.EXPECT().Ping(gomock.Any(), client.PingOptions{NegotiateAPIVersion: true}). + Return(client.PingResult{}, errors.New("engine unreachable")).AnyTimes() + // The dialer must never be reached when negotiation fails. + apiClient.EXPECT().Dialer().Times(0) + + tested, err := NewComposeService(cli) + assert.NilError(t, err) + s := tested.(*composeService) + + err = s.pushProjectConfig(t.Context(), newUpTestProject(), true) + assert.ErrorContains(t, err, "negotiating API version") + assert.ErrorContains(t, err, "engine unreachable") +} + +func TestHasActiveProfiles(t *testing.T) { + tests := []struct { + name string + profiles []string + want bool + }{ + {name: "nil", profiles: nil, want: false}, + {name: "empty slice", profiles: []string{}, want: false}, + {name: "single blank (default, no profile)", profiles: []string{""}, want: false}, + {name: "named profile", profiles: []string{"debug"}, want: true}, + {name: "blank plus named", profiles: []string{"", "debug"}, want: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, hasActiveProfiles(tt.profiles), tt.want) + }) + } +} + func TestShouldFollowStartEvent(t *testing.T) { tests := []struct { name string