From 1542f5e969e7793e69c687838f486cb438f629e5 Mon Sep 17 00:00:00 2001 From: Dennis Chen Date: Fri, 1 Jun 2018 15:21:29 +0800 Subject: [PATCH 01/13] Refactor and cleanup the intermediate container creation This PR is trying to refactor the `probeAndCreate` and cleanup related codes based on the refactoring. Signed-off-by: Dennis Chen Upstream-commit: 7f280f6f6514598a31d2523d4b9437ca07ee0526 Component: engine --- .../engine/builder/dockerfile/dispatchers.go | 8 ++++---- components/engine/builder/dockerfile/internals.go | 15 +++++---------- 2 files changed, 9 insertions(+), 14 deletions(-) diff --git a/components/engine/builder/dockerfile/dispatchers.go b/components/engine/builder/dockerfile/dispatchers.go index 9dd7502453..8382fa8991 100644 --- a/components/engine/builder/dockerfile/dispatchers.go +++ b/components/engine/builder/dockerfile/dispatchers.go @@ -28,7 +28,6 @@ import ( "github.com/docker/docker/pkg/system" "github.com/docker/go-connections/nat" "github.com/pkg/errors" - "github.com/sirupsen/logrus" ) // ENV foo bar @@ -305,10 +304,12 @@ func dispatchWorkdir(d dispatchRequest, c *instructions.WorkdirCommand) error { comment := "WORKDIR " + runConfig.WorkingDir runConfigWithCommentCmd := copyRunConfig(runConfig, withCmdCommentString(comment, d.state.operatingSystem)) + containerID, err := d.builder.probeAndCreate(d.state, runConfigWithCommentCmd) if err != nil || containerID == "" { return err } + if err := d.builder.docker.ContainerCreateWorkdir(containerID); err != nil { return err } @@ -350,8 +351,7 @@ func dispatchRun(d dispatchRequest, c *instructions.RunCommand) error { runConfigForCacheProbe := copyRunConfig(stateRunConfig, withCmd(saveCmd), withEntrypointOverride(saveCmd, nil)) - hit, err := d.builder.probeCache(d.state, runConfigForCacheProbe) - if err != nil || hit { + if hit, err := d.builder.probeCache(d.state, runConfigForCacheProbe); err != nil || hit { return err } @@ -363,11 +363,11 @@ func dispatchRun(d dispatchRequest, c *instructions.RunCommand) error { // set config as already being escaped, this prevents double escaping on windows runConfig.ArgsEscaped = true - logrus.Debugf("[BUILDER] Command to be executed: %v", runConfig.Cmd) cID, err := d.builder.create(runConfig) if err != nil { return err } + if err := d.builder.containerManager.Run(d.builder.clientCtx, cID, d.builder.Stdout, d.builder.Stderr); err != nil { if err, ok := err.(*statusCodeError); ok { // TODO: change error type, because jsonmessage.JSONError assumes HTTP diff --git a/components/engine/builder/dockerfile/internals.go b/components/engine/builder/dockerfile/internals.go index 53748f0619..88e75a2179 100644 --- a/components/engine/builder/dockerfile/internals.go +++ b/components/engine/builder/dockerfile/internals.go @@ -27,6 +27,7 @@ import ( "github.com/docker/docker/pkg/system" "github.com/docker/go-connections/nat" "github.com/pkg/errors" + "github.com/sirupsen/logrus" ) // Archiver defines an interface for copying files from one destination to @@ -84,12 +85,8 @@ func (b *Builder) commit(dispatchState *dispatchState, comment string) error { } runConfigWithCommentCmd := copyRunConfig(dispatchState.runConfig, withCmdComment(comment, dispatchState.operatingSystem)) - hit, err := b.probeCache(dispatchState, runConfigWithCommentCmd) - if err != nil || hit { - return err - } - id, err := b.create(runConfigWithCommentCmd) - if err != nil { + id, err := b.probeAndCreate(dispatchState, runConfigWithCommentCmd) + if err != nil || id == "" { return err } @@ -413,13 +410,11 @@ func (b *Builder) probeAndCreate(dispatchState *dispatchState, runConfig *contai if hit, err := b.probeCache(dispatchState, runConfig); err != nil || hit { return "", err } - // Set a log config to override any default value set on the daemon - hostConfig := &container.HostConfig{LogConfig: defaultLogConfig} - container, err := b.containerManager.Create(runConfig, hostConfig) - return container.ID, err + return b.create(runConfig) } func (b *Builder) create(runConfig *container.Config) (string, error) { + logrus.Debugf("[BUILDER] Command to be executed: %v", runConfig.Cmd) hostConfig := hostConfigFromOptions(b.options) container, err := b.containerManager.Create(runConfig, hostConfig) if err != nil { From 8a3da41a28ec50386216822b27a7ab0ad309d11b Mon Sep 17 00:00:00 2001 From: Arash Deshmeh Date: Sat, 2 Jun 2018 09:44:45 -0400 Subject: [PATCH 02/13] Added network package to integration/internal to refactor integration tests calls to client.NetworkCreate Signed-off-by: Arash Deshmeh Upstream-commit: 0e3012cfff1e81f7c779b1310f6069dc909d7987 Component: engine --- .../integration/internal/network/network.go | 35 ++++++++++++ .../integration/internal/network/ops.go | 57 +++++++++++++++++++ 2 files changed, 92 insertions(+) create mode 100644 components/engine/integration/internal/network/network.go create mode 100644 components/engine/integration/internal/network/ops.go diff --git a/components/engine/integration/internal/network/network.go b/components/engine/integration/internal/network/network.go new file mode 100644 index 0000000000..b9550362f0 --- /dev/null +++ b/components/engine/integration/internal/network/network.go @@ -0,0 +1,35 @@ +package network + +import ( + "context" + "testing" + + "github.com/docker/docker/api/types" + "github.com/docker/docker/client" + "github.com/gotestyourself/gotestyourself/assert" +) + +func createNetwork(ctx context.Context, client client.APIClient, name string, ops ...func(*types.NetworkCreate)) (string, error) { + config := types.NetworkCreate{} + + for _, op := range ops { + op(&config) + } + + n, err := client.NetworkCreate(ctx, name, config) + return n.ID, err +} + +// Create creates a network with the specified options +func Create(ctx context.Context, client client.APIClient, name string, ops ...func(*types.NetworkCreate)) (string, error) { + return createNetwork(ctx, client, name, ops...) +} + +// CreateNoError creates a network with the specified options and verifies there were no errors +func CreateNoError(t *testing.T, ctx context.Context, client client.APIClient, name string, ops ...func(*types.NetworkCreate)) string { // nolint: golint + t.Helper() + + name, err := createNetwork(ctx, client, name, ops...) + assert.NilError(t, err) + return name +} diff --git a/components/engine/integration/internal/network/ops.go b/components/engine/integration/internal/network/ops.go new file mode 100644 index 0000000000..f7639ff307 --- /dev/null +++ b/components/engine/integration/internal/network/ops.go @@ -0,0 +1,57 @@ +package network + +import ( + "github.com/docker/docker/api/types" + "github.com/docker/docker/api/types/network" +) + +// WithDriver sets the driver of the network +func WithDriver(driver string) func(*types.NetworkCreate) { + return func(n *types.NetworkCreate) { + n.Driver = driver + } +} + +// WithIPv6 Enables IPv6 on the network +func WithIPv6() func(*types.NetworkCreate) { + return func(n *types.NetworkCreate) { + n.EnableIPv6 = true + } +} + +// WithMacvlan sets the network as macvlan with the specified parent +func WithMacvlan(parent string) func(*types.NetworkCreate) { + return func(n *types.NetworkCreate) { + n.Driver = "macvlan" + if parent != "" { + n.Options = map[string]string{ + "parent": parent, + } + } + } +} + +// WithOption adds the specified key/value pair to network's options +func WithOption(key, value string) func(*types.NetworkCreate) { + return func(n *types.NetworkCreate) { + if n.Options == nil { + n.Options = map[string]string{} + } + n.Options[key] = value + } +} + +// WithIPAM adds an IPAM with the specified Subnet and Gateway to the network +func WithIPAM(subnet, gateway string) func(*types.NetworkCreate) { + return func(n *types.NetworkCreate) { + if n.IPAM == nil { + n.IPAM = &network.IPAM{} + } + + n.IPAM.Config = append(n.IPAM.Config, network.IPAMConfig{ + Subnet: subnet, + Gateway: gateway, + AuxAddress: map[string]string{}, + }) + } +} From f422827dd45a81e7fde2aaa987b540a25fc0bf14 Mon Sep 17 00:00:00 2001 From: Arash Deshmeh Date: Sat, 2 Jun 2018 09:45:59 -0400 Subject: [PATCH 03/13] refactored integration tests under integration/network/macvlan to use network.Create Signed-off-by: Arash Deshmeh Upstream-commit: 0418893f0b5a7413edd950cfabd2337161edc705 Component: engine --- .../network/macvlan/macvlan_test.go | 139 +++++++----------- 1 file changed, 50 insertions(+), 89 deletions(-) diff --git a/components/engine/integration/network/macvlan/macvlan_test.go b/components/engine/integration/network/macvlan/macvlan_test.go index 7358af873b..7d9acd813e 100644 --- a/components/engine/integration/network/macvlan/macvlan_test.go +++ b/components/engine/integration/network/macvlan/macvlan_test.go @@ -7,9 +7,9 @@ import ( "time" "github.com/docker/docker/api/types" - "github.com/docker/docker/api/types/network" "github.com/docker/docker/client" "github.com/docker/docker/integration/internal/container" + net "github.com/docker/docker/integration/internal/network" n "github.com/docker/docker/integration/network" "github.com/docker/docker/internal/test/daemon" "github.com/gotestyourself/gotestyourself/assert" @@ -33,16 +33,13 @@ func TestDockerNetworkMacvlanPersistance(t *testing.T) { client, err := d.NewClient() assert.NilError(t, err) - _, err = client.NetworkCreate(context.Background(), "dm-persist", types.NetworkCreate{ - Driver: "macvlan", - Options: map[string]string{ - "parent": "dm-dummy0.60", - }, - }) - assert.NilError(t, err) - assert.Check(t, n.IsNetworkAvailable(client, "dm-persist")) + netName := "dm-persist" + net.CreateNoError(t, context.Background(), client, netName, + net.WithMacvlan("dm-dummy0.60"), + ) + assert.Check(t, n.IsNetworkAvailable(client, netName)) d.Restart(t) - assert.Check(t, n.IsNetworkAvailable(client, "dm-persist")) + assert.Check(t, n.IsNetworkAvailable(client, netName)) } func TestDockerNetworkMacvlan(t *testing.T) { @@ -91,29 +88,25 @@ func testMacvlanOverlapParent(client client.APIClient) func(*testing.T) { n.CreateMasterDummy(t, master) defer n.DeleteInterface(t, master) - _, err := client.NetworkCreate(context.Background(), "dm-subinterface", types.NetworkCreate{ - Driver: "macvlan", - Options: map[string]string{ - "parent": "dm-dummy0.40", - }, - }) - assert.NilError(t, err) - assert.Check(t, n.IsNetworkAvailable(client, "dm-subinterface")) + netName := "dm-subinterface" + parentName := "dm-dummy0.40" + net.CreateNoError(t, context.Background(), client, netName, + net.WithMacvlan(parentName), + ) + assert.Check(t, n.IsNetworkAvailable(client, netName)) - _, err = client.NetworkCreate(context.Background(), "dm-parent-net-overlap", types.NetworkCreate{ - Driver: "macvlan", - Options: map[string]string{ - "parent": "dm-dummy0.40", - }, - }) + _, err := net.Create(context.Background(), client, "dm-parent-net-overlap", + net.WithMacvlan(parentName), + ) assert.Check(t, err != nil) + // delete the network while preserving the parent link - err = client.NetworkRemove(context.Background(), "dm-subinterface") + err = client.NetworkRemove(context.Background(), netName) assert.NilError(t, err) - assert.Check(t, n.IsNetworkNotAvailable(client, "dm-subinterface")) + assert.Check(t, n.IsNetworkNotAvailable(client, netName)) // verify the network delete did not delete the predefined link - n.LinkExists(t, "dm-dummy0") + n.LinkExists(t, master) } } @@ -121,26 +114,24 @@ func testMacvlanSubinterface(client client.APIClient) func(*testing.T) { return func(t *testing.T) { // verify the same parent interface cannot be used if already in use by an existing network master := "dm-dummy0" + parentName := "dm-dummy0.20" n.CreateMasterDummy(t, master) defer n.DeleteInterface(t, master) - n.CreateVlanInterface(t, master, "dm-dummy0.20", "20") + n.CreateVlanInterface(t, master, parentName, "20") - _, err := client.NetworkCreate(context.Background(), "dm-subinterface", types.NetworkCreate{ - Driver: "macvlan", - Options: map[string]string{ - "parent": "dm-dummy0.20", - }, - }) - assert.NilError(t, err) - assert.Check(t, n.IsNetworkAvailable(client, "dm-subinterface")) + netName := "dm-subinterface" + net.CreateNoError(t, context.Background(), client, netName, + net.WithMacvlan(parentName), + ) + assert.Check(t, n.IsNetworkAvailable(client, netName)) // delete the network while preserving the parent link - err = client.NetworkRemove(context.Background(), "dm-subinterface") + err := client.NetworkRemove(context.Background(), netName) assert.NilError(t, err) - assert.Check(t, n.IsNetworkNotAvailable(client, "dm-subinterface")) + assert.Check(t, n.IsNetworkNotAvailable(client, netName)) // verify the network delete did not delete the predefined link - n.LinkExists(t, "dm-dummy0.20") + n.LinkExists(t, parentName) } } @@ -190,34 +181,17 @@ func testMacvlanInternalMode(client client.APIClient) func(*testing.T) { func testMacvlanMultiSubnet(client client.APIClient) func(*testing.T) { return func(t *testing.T) { - _, err := client.NetworkCreate(context.Background(), "dualstackbridge", types.NetworkCreate{ - Driver: "macvlan", - EnableIPv6: true, - IPAM: &network.IPAM{ - Config: []network.IPAMConfig{ - { - Subnet: "172.28.100.0/24", - AuxAddress: map[string]string{}, - }, - { - Subnet: "172.28.102.0/24", - Gateway: "172.28.102.254", - AuxAddress: map[string]string{}, - }, - { - Subnet: "2001:db8:abc2::/64", - AuxAddress: map[string]string{}, - }, - { - Subnet: "2001:db8:abc4::/64", - Gateway: "2001:db8:abc4::254", - AuxAddress: map[string]string{}, - }, - }, - }, - }) - assert.NilError(t, err) - assert.Check(t, n.IsNetworkAvailable(client, "dualstackbridge")) + netName := "dualstackbridge" + net.CreateNoError(t, context.Background(), client, netName, + net.WithMacvlan(""), + net.WithIPv6(), + net.WithIPAM("172.28.100.0/24", ""), + net.WithIPAM("172.28.102.0/24", "172.28.102.254"), + net.WithIPAM("2001:db8:abc2::/64", ""), + net.WithIPAM("2001:db8:abc4::/64", "2001:db8:abc4::254"), + ) + + assert.Check(t, n.IsNetworkAvailable(client, netName)) // start dual stack containers and verify the user specified --ip and --ip6 addresses on subnets 172.28.100.0/24 and 2001:db8:abc2::/64 ctx := context.Background() @@ -276,28 +250,15 @@ func testMacvlanMultiSubnet(client client.APIClient) func(*testing.T) { func testMacvlanAddressing(client client.APIClient) func(*testing.T) { return func(t *testing.T) { // Ensure the default gateways, next-hops and default dev devices are properly set - _, err := client.NetworkCreate(context.Background(), "dualstackbridge", types.NetworkCreate{ - Driver: "macvlan", - EnableIPv6: true, - Options: map[string]string{ - "macvlan_mode": "bridge", - }, - IPAM: &network.IPAM{ - Config: []network.IPAMConfig{ - { - Subnet: "172.28.130.0/24", - AuxAddress: map[string]string{}, - }, - { - Subnet: "2001:db8:abca::/64", - Gateway: "2001:db8:abca::254", - AuxAddress: map[string]string{}, - }, - }, - }, - }) - assert.NilError(t, err) - assert.Check(t, n.IsNetworkAvailable(client, "dualstackbridge")) + netName := "dualstackbridge" + net.CreateNoError(t, context.Background(), client, netName, + net.WithMacvlan(""), + net.WithIPv6(), + net.WithOption("macvlan_mode", "bridge"), + net.WithIPAM("172.28.130.0/24", ""), + net.WithIPAM("2001:db8:abca::/64", "2001:db8:abca::254"), + ) + assert.Check(t, n.IsNetworkAvailable(client, netName)) ctx := context.Background() id1 := container.Run(t, ctx, client, From 9022898c18d39064c987ec588cd434caf02145b5 Mon Sep 17 00:00:00 2001 From: Arash Deshmeh Date: Tue, 5 Jun 2018 15:50:46 -0400 Subject: [PATCH 04/13] integration tests under integration/container/links_linux_test.go use unique names Signed-off-by: Arash Deshmeh Upstream-commit: 077247050df0bf97b7b62e439152eb362e9ec813 Component: engine --- .../engine/integration/container/links_linux_test.go | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/components/engine/integration/container/links_linux_test.go b/components/engine/integration/container/links_linux_test.go index 4fafd32c85..9baa32728d 100644 --- a/components/engine/integration/container/links_linux_test.go +++ b/components/engine/integration/container/links_linux_test.go @@ -41,15 +41,17 @@ func TestLinksContainerNames(t *testing.T) { client := request.NewAPIClient(t) ctx := context.Background() - container.Run(t, ctx, client, container.WithName("first")) - container.Run(t, ctx, client, container.WithName("second"), container.WithLinks("first:first")) + containerA := "first_" + t.Name() + containerB := "second_" + t.Name() + container.Run(t, ctx, client, container.WithName(containerA)) + container.Run(t, ctx, client, container.WithName(containerB), container.WithLinks(containerA+":"+containerA)) - f := filters.NewArgs(filters.Arg("name", "first")) + f := filters.NewArgs(filters.Arg("name", containerA)) containers, err := client.ContainerList(ctx, types.ContainerListOptions{ Filters: f, }) assert.NilError(t, err) assert.Check(t, is.Equal(1, len(containers))) - assert.Check(t, is.DeepEqual([]string{"/first", "/second/first"}, containers[0].Names)) + assert.Check(t, is.DeepEqual([]string{"/" + containerA, "/" + containerB + "/" + containerA}, containers[0].Names)) } From 43d60463c28796d335e792dc86dc3f2e22b37b28 Mon Sep 17 00:00:00 2001 From: Vincent Demeester Date: Fri, 1 Jun 2018 12:47:38 +0200 Subject: [PATCH 05/13] Add support for `init` on services It's already supported by `swarmkit`, and act the same as `HostConfig.Init` on container creation. Signed-off-by: Vincent Demeester Upstream-commit: e401b88e59e098745744917c555d549f08353e6d Component: engine --- components/engine/api/swagger.yaml | 4 ++ .../engine/api/types/swarm/container.go | 1 + .../daemon/cluster/convert/container.go | 17 ++++++ .../cluster/executor/container/container.go | 9 ++++ .../integration/internal/swarm/service.go | 8 +++ .../engine/integration/service/create_test.go | 54 +++++++++++++++++++ .../engine/internal/test/daemon/daemon.go | 6 ++- components/engine/internal/test/daemon/ops.go | 6 +++ 8 files changed, 104 insertions(+), 1 deletion(-) diff --git a/components/engine/api/swagger.yaml b/components/engine/api/swagger.yaml index 3cdfe2988c..86374415fd 100644 --- a/components/engine/api/swagger.yaml +++ b/components/engine/api/swagger.yaml @@ -2721,6 +2721,10 @@ definitions: - "default" - "process" - "hyperv" + Init: + description: "Run an init inside the container that forwards signals and reaps processes. This field is omitted if empty, and the default (as configured on the daemon) is used." + type: "boolean" + x-nullable: true NetworkAttachmentSpec: description: | Read-only spec type for non-swarm containers attached to swarm overlay diff --git a/components/engine/api/types/swarm/container.go b/components/engine/api/types/swarm/container.go index 0041653c9d..151211ff5a 100644 --- a/components/engine/api/types/swarm/container.go +++ b/components/engine/api/types/swarm/container.go @@ -55,6 +55,7 @@ type ContainerSpec struct { User string `json:",omitempty"` Groups []string `json:",omitempty"` Privileges *Privileges `json:",omitempty"` + Init *bool `json:",omitempty"` StopSignal string `json:",omitempty"` TTY bool `json:",omitempty"` OpenStdin bool `json:",omitempty"` diff --git a/components/engine/daemon/cluster/convert/container.go b/components/engine/daemon/cluster/convert/container.go index 0a34fc73e4..d889b4004c 100644 --- a/components/engine/daemon/cluster/convert/container.go +++ b/components/engine/daemon/cluster/convert/container.go @@ -35,6 +35,7 @@ func containerSpecFromGRPC(c *swarmapi.ContainerSpec) *types.ContainerSpec { Secrets: secretReferencesFromGRPC(c.Secrets), Configs: configReferencesFromGRPC(c.Configs), Isolation: IsolationFromGRPC(c.Isolation), + Init: initFromGRPC(c.Init), } if c.DNSConfig != nil { @@ -119,6 +120,21 @@ func containerSpecFromGRPC(c *swarmapi.ContainerSpec) *types.ContainerSpec { return containerSpec } +func initFromGRPC(v *gogotypes.BoolValue) *bool { + if v == nil { + return nil + } + value := v.GetValue() + return &value +} + +func initToGRPC(v *bool) *gogotypes.BoolValue { + if v == nil { + return nil + } + return &gogotypes.BoolValue{Value: *v} +} + func secretReferencesToGRPC(sr []*types.SecretReference) []*swarmapi.SecretReference { refs := make([]*swarmapi.SecretReference, 0, len(sr)) for _, s := range sr { @@ -234,6 +250,7 @@ func containerToGRPC(c *types.ContainerSpec) (*swarmapi.ContainerSpec, error) { Secrets: secretReferencesToGRPC(c.Secrets), Configs: configReferencesToGRPC(c.Configs), Isolation: isolationToGRPC(c.Isolation), + Init: initToGRPC(c.Init), } if c.DNSConfig != nil { diff --git a/components/engine/daemon/cluster/executor/container/container.go b/components/engine/daemon/cluster/executor/container/container.go index 69d673bd30..77d21d2c1f 100644 --- a/components/engine/daemon/cluster/executor/container/container.go +++ b/components/engine/daemon/cluster/executor/container/container.go @@ -172,6 +172,14 @@ func (c *containerConfig) isolation() enginecontainer.Isolation { return convert.IsolationFromGRPC(c.spec().Isolation) } +func (c *containerConfig) init() *bool { + if c.spec().Init == nil { + return nil + } + init := c.spec().Init.GetValue() + return &init +} + func (c *containerConfig) exposedPorts() map[nat.Port]struct{} { exposedPorts := make(map[nat.Port]struct{}) if c.task.Endpoint == nil { @@ -355,6 +363,7 @@ func (c *containerConfig) hostConfig() *enginecontainer.HostConfig { Mounts: c.mounts(), ReadonlyRootfs: c.spec().ReadOnly, Isolation: c.isolation(), + Init: c.init(), } if c.spec().DNSConfig != nil { diff --git a/components/engine/integration/internal/swarm/service.go b/components/engine/integration/internal/swarm/service.go index e6a1bfcdd0..5567ad6ede 100644 --- a/components/engine/integration/internal/swarm/service.go +++ b/components/engine/integration/internal/swarm/service.go @@ -86,6 +86,14 @@ func defaultServiceSpec() swarmtypes.ServiceSpec { return spec } +// ServiceWithInit sets whether the service should use init or not +func ServiceWithInit(b *bool) func(*swarmtypes.ServiceSpec) { + return func(spec *swarmtypes.ServiceSpec) { + ensureContainerSpec(spec) + spec.TaskTemplate.ContainerSpec.Init = b + } +} + // ServiceWithImage sets the image to use for the service func ServiceWithImage(image string) func(*swarmtypes.ServiceSpec) { return func(spec *swarmtypes.ServiceSpec) { diff --git a/components/engine/integration/service/create_test.go b/components/engine/integration/service/create_test.go index 68af6e825a..517ee0d514 100644 --- a/components/engine/integration/service/create_test.go +++ b/components/engine/integration/service/create_test.go @@ -2,6 +2,7 @@ package service // import "github.com/docker/docker/integration/service" import ( "context" + "fmt" "io/ioutil" "testing" "time" @@ -11,11 +12,64 @@ import ( swarmtypes "github.com/docker/docker/api/types/swarm" "github.com/docker/docker/client" "github.com/docker/docker/integration/internal/swarm" + "github.com/docker/docker/internal/test/daemon" "github.com/gotestyourself/gotestyourself/assert" is "github.com/gotestyourself/gotestyourself/assert/cmp" "github.com/gotestyourself/gotestyourself/poll" ) +func TestServiceCreateInit(t *testing.T) { + defer setupTest(t)() + t.Run("daemonInitDisabled", testServiceCreateInit(false)) + t.Run("daemonInitEnabled", testServiceCreateInit(true)) +} + +func testServiceCreateInit(daemonEnabled bool) func(t *testing.T) { + return func(t *testing.T) { + var ops = []func(*daemon.Daemon){} + + if daemonEnabled { + ops = append(ops, daemon.WithInit) + } + d := swarm.NewSwarm(t, testEnv, ops...) + defer d.Stop(t) + client := d.NewClientT(t) + defer client.Close() + + booleanTrue := true + booleanFalse := false + + serviceID := swarm.CreateService(t, d) + poll.WaitOn(t, serviceRunningTasksCount(client, serviceID, 1), swarm.ServicePoll) + i := inspectServiceContainer(t, client, serviceID) + // HostConfig.Init == nil means that it delegates to daemon configuration + assert.Check(t, i.HostConfig.Init == nil) + + serviceID = swarm.CreateService(t, d, swarm.ServiceWithInit(&booleanTrue)) + poll.WaitOn(t, serviceRunningTasksCount(client, serviceID, 1), swarm.ServicePoll) + i = inspectServiceContainer(t, client, serviceID) + assert.Check(t, is.Equal(true, *i.HostConfig.Init)) + + serviceID = swarm.CreateService(t, d, swarm.ServiceWithInit(&booleanFalse)) + poll.WaitOn(t, serviceRunningTasksCount(client, serviceID, 1), swarm.ServicePoll) + i = inspectServiceContainer(t, client, serviceID) + assert.Check(t, is.Equal(false, *i.HostConfig.Init)) + } +} + +func inspectServiceContainer(t *testing.T, client client.APIClient, serviceID string) types.ContainerJSON { + t.Helper() + filter := filters.NewArgs() + filter.Add("label", fmt.Sprintf("com.docker.swarm.service.id=%s", serviceID)) + containers, err := client.ContainerList(context.Background(), types.ContainerListOptions{Filters: filter}) + assert.NilError(t, err) + assert.Check(t, is.Len(containers, 1)) + + i, err := client.ContainerInspect(context.Background(), containers[0].ID) + assert.NilError(t, err) + return i +} + func TestCreateServiceMultipleTimes(t *testing.T) { defer setupTest(t)() d := swarm.NewSwarm(t, testEnv) diff --git a/components/engine/internal/test/daemon/daemon.go b/components/engine/internal/test/daemon/daemon.go index 9ba13edc0a..a0d7ed4855 100644 --- a/components/engine/internal/test/daemon/daemon.go +++ b/components/engine/internal/test/daemon/daemon.go @@ -66,6 +66,7 @@ type Daemon struct { userlandProxy bool execRoot string experimental bool + init bool dockerdBinary string log logT @@ -229,7 +230,10 @@ func (d *Daemon) StartWithLogFile(out *os.File, providedArgs ...string) error { fmt.Sprintf("--userland-proxy=%t", d.userlandProxy), ) if d.experimental { - args = append(args, "--experimental", "--init") + args = append(args, "--experimental") + } + if d.init { + args = append(args, "--init") } if !(d.UseDefaultHost || d.UseDefaultTLSHost) { args = append(args, []string{"--host", d.Sock()}...) diff --git a/components/engine/internal/test/daemon/ops.go b/components/engine/internal/test/daemon/ops.go index 288fe88070..34db073b57 100644 --- a/components/engine/internal/test/daemon/ops.go +++ b/components/engine/internal/test/daemon/ops.go @@ -5,6 +5,12 @@ import "github.com/docker/docker/internal/test/environment" // WithExperimental sets the daemon in experimental mode func WithExperimental(d *Daemon) { d.experimental = true + d.init = true +} + +// WithInit sets the daemon init +func WithInit(d *Daemon) { + d.init = true } // WithDockerdBinary sets the dockerd binary to the specified one From 29579f20f8f7d5bc612cf0b29270b21a1cef1cfb Mon Sep 17 00:00:00 2001 From: Arash Deshmeh Date: Thu, 7 Jun 2018 11:55:08 -0400 Subject: [PATCH 06/13] use unique names for resources in create service integration tests Signed-off-by: Arash Deshmeh Upstream-commit: d5ae23fc11b91275801dcab2326e803725a645a6 Component: engine --- .../engine/integration/service/create_test.go | 30 +++++++++++-------- 1 file changed, 18 insertions(+), 12 deletions(-) diff --git a/components/engine/integration/service/create_test.go b/components/engine/integration/service/create_test.go index 68af6e825a..e37b5ed5e3 100644 --- a/components/engine/integration/service/create_test.go +++ b/components/engine/integration/service/create_test.go @@ -23,7 +23,7 @@ func TestCreateServiceMultipleTimes(t *testing.T) { client := d.NewClientT(t) defer client.Close() - overlayName := "overlay1" + overlayName := "overlay1_" + t.Name() networkCreate := types.NetworkCreate{ CheckDuplicate: true, Driver: "overlay", @@ -35,9 +35,10 @@ func TestCreateServiceMultipleTimes(t *testing.T) { var instances uint64 = 4 + serviceName := "TestService_" + t.Name() serviceSpec := []swarm.ServiceSpecOpt{ swarm.ServiceWithReplicas(instances), - swarm.ServiceWithName("TestService"), + swarm.ServiceWithName(serviceName), swarm.ServiceWithNetwork(overlayName), } @@ -75,7 +76,7 @@ func TestCreateWithDuplicateNetworkNames(t *testing.T) { client := d.NewClientT(t) defer client.Close() - name := "foo" + name := "foo_" + t.Name() networkCreate := types.NetworkCreate{ CheckDuplicate: false, Driver: "bridge", @@ -95,9 +96,10 @@ func TestCreateWithDuplicateNetworkNames(t *testing.T) { // Create Service with the same name var instances uint64 = 1 + serviceName := "top_" + t.Name() serviceID := swarm.CreateService(t, d, swarm.ServiceWithReplicas(instances), - swarm.ServiceWithName("top"), + swarm.ServiceWithName(serviceName), swarm.ServiceWithNetwork(name), ) @@ -138,18 +140,20 @@ func TestCreateServiceSecretFileMode(t *testing.T) { defer client.Close() ctx := context.Background() + secretName := "TestSecret_" + t.Name() secretResp, err := client.SecretCreate(ctx, swarmtypes.SecretSpec{ Annotations: swarmtypes.Annotations{ - Name: "TestSecret", + Name: secretName, }, Data: []byte("TESTSECRET"), }) assert.NilError(t, err) var instances uint64 = 1 + serviceName := "TestService_" + t.Name() serviceID := swarm.CreateService(t, d, swarm.ServiceWithReplicas(instances), - swarm.ServiceWithName("TestService"), + swarm.ServiceWithName(serviceName), swarm.ServiceWithCommand([]string{"/bin/sh", "-c", "ls -l /etc/secret || /bin/top"}), swarm.ServiceWithSecret(&swarmtypes.SecretReference{ File: &swarmtypes.SecretReferenceFileTarget{ @@ -159,7 +163,7 @@ func TestCreateServiceSecretFileMode(t *testing.T) { Mode: 0777, }, SecretID: secretResp.ID, - SecretName: "TestSecret", + SecretName: secretName, }), ) @@ -189,7 +193,7 @@ func TestCreateServiceSecretFileMode(t *testing.T) { poll.WaitOn(t, serviceIsRemoved(client, serviceID), swarm.ServicePoll) poll.WaitOn(t, noTasks(client), swarm.ServicePoll) - err = client.SecretRemove(ctx, "TestSecret") + err = client.SecretRemove(ctx, secretName) assert.NilError(t, err) } @@ -201,17 +205,19 @@ func TestCreateServiceConfigFileMode(t *testing.T) { defer client.Close() ctx := context.Background() + configName := "TestConfig_" + t.Name() configResp, err := client.ConfigCreate(ctx, swarmtypes.ConfigSpec{ Annotations: swarmtypes.Annotations{ - Name: "TestConfig", + Name: configName, }, Data: []byte("TESTCONFIG"), }) assert.NilError(t, err) var instances uint64 = 1 + serviceName := "TestService_" + t.Name() serviceID := swarm.CreateService(t, d, - swarm.ServiceWithName("TestService"), + swarm.ServiceWithName(serviceName), swarm.ServiceWithCommand([]string{"/bin/sh", "-c", "ls -l /etc/config || /bin/top"}), swarm.ServiceWithReplicas(instances), swarm.ServiceWithConfig(&swarmtypes.ConfigReference{ @@ -222,7 +228,7 @@ func TestCreateServiceConfigFileMode(t *testing.T) { Mode: 0777, }, ConfigID: configResp.ID, - ConfigName: "TestConfig", + ConfigName: configName, }), ) @@ -252,7 +258,7 @@ func TestCreateServiceConfigFileMode(t *testing.T) { poll.WaitOn(t, serviceIsRemoved(client, serviceID)) poll.WaitOn(t, noTasks(client)) - err = client.ConfigRemove(ctx, "TestConfig") + err = client.ConfigRemove(ctx, configName) assert.NilError(t, err) } From c2a60aeb4580639ca59f88d388ed191d53729162 Mon Sep 17 00:00:00 2001 From: Vincent Demeester Date: Thu, 7 Jun 2018 19:50:41 +0200 Subject: [PATCH 07/13] =?UTF-8?q?Mark=20@thajeztah=20as=20a=20MAINTAINER?= =?UTF-8?q?=E2=80=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit … discovered recently that our very own Sebastiaan was not even listed in the maintainers, so fixing this mistake ! Signed-off-by: Vincent Demeester Upstream-commit: 7f9dd507d3f50cd343480cb0776a3ccca83c5764 Component: engine --- components/engine/MAINTAINERS | 1 + 1 file changed, 1 insertion(+) diff --git a/components/engine/MAINTAINERS b/components/engine/MAINTAINERS index 8f76e71010..3ac06d2728 100644 --- a/components/engine/MAINTAINERS +++ b/components/engine/MAINTAINERS @@ -40,6 +40,7 @@ "mlaventure", "runcom", "stevvooe", + "thajeztah", "tianon", "tibor", "tonistiigi", From dc6063672cd4f9edcc3b6f8652166da965218e74 Mon Sep 17 00:00:00 2001 From: Joe Ferguson Date: Fri, 21 Jul 2017 17:47:15 -0700 Subject: [PATCH 08/13] Close readclosers returned by DecompressStream Signed-off-by: Joe Ferguson Upstream-commit: 76e99e1a8b60c831415d2f6a6a7954e16c25620b Component: engine --- components/engine/pkg/archive/archive.go | 1 + components/engine/pkg/archive/diff.go | 4 +++- components/engine/plugin/blobstore.go | 1 + 3 files changed, 5 insertions(+), 1 deletion(-) diff --git a/components/engine/pkg/archive/archive.go b/components/engine/pkg/archive/archive.go index 5fb3995e9b..daddebded4 100644 --- a/components/engine/pkg/archive/archive.go +++ b/components/engine/pkg/archive/archive.go @@ -127,6 +127,7 @@ func IsArchivePath(path string) bool { if err != nil { return false } + defer rdr.Close() r := tar.NewReader(rdr) _, err = r.Next() return err == nil diff --git a/components/engine/pkg/archive/diff.go b/components/engine/pkg/archive/diff.go index d0cff98ffc..fae4b9de02 100644 --- a/components/engine/pkg/archive/diff.go +++ b/components/engine/pkg/archive/diff.go @@ -247,10 +247,12 @@ func applyLayerHandler(dest string, layer io.Reader, options *TarOptions, decomp defer system.Umask(oldmask) // ignore err, ErrNotSupportedPlatform if decompress { - layer, err = DecompressStream(layer) + decompLayer, err := DecompressStream(layer) if err != nil { return 0, err } + defer decompLayer.Close() + layer = decompLayer } return UnpackLayer(dest, layer, options) } diff --git a/components/engine/plugin/blobstore.go b/components/engine/plugin/blobstore.go index fd7f040efc..a24e7bdf4f 100644 --- a/components/engine/plugin/blobstore.go +++ b/components/engine/plugin/blobstore.go @@ -145,6 +145,7 @@ func (dm *downloadManager) Download(ctx context.Context, initialRootFS image.Roo if err != nil { return initialRootFS, nil, err } + defer inflatedLayerData.Close() digester := digest.Canonical.Digester() if _, err := chrootarchive.ApplyLayer(dm.tmpDir, io.TeeReader(inflatedLayerData, digester.Hash())); err != nil { return initialRootFS, nil, err From 97950284c352ec4927309b17cf5d80fce3cd6eb8 Mon Sep 17 00:00:00 2001 From: Daniel Nephin Date: Thu, 7 Jun 2018 18:26:07 -0400 Subject: [PATCH 09/13] Add image metrics for push and pull Signed-off-by: Daniel Nephin Upstream-commit: 6910019bbebd71d6bb5c949a40e96b49a6b41f45 Component: engine --- components/engine/daemon/images/image_pull.go | 6 +++++- components/engine/daemon/images/image_push.go | 3 +++ 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/components/engine/daemon/images/image_pull.go b/components/engine/daemon/images/image_pull.go index e25e1d0631..238c38b6b3 100644 --- a/components/engine/daemon/images/image_pull.go +++ b/components/engine/daemon/images/image_pull.go @@ -5,6 +5,7 @@ import ( "io" "runtime" "strings" + "time" dist "github.com/docker/distribution" "github.com/docker/distribution/reference" @@ -20,6 +21,7 @@ import ( // PullImage initiates a pull operation. image is the repository name to pull, and // tag may be either empty, or indicate a specific tag to pull. func (i *ImageService) PullImage(ctx context.Context, image, tag, os string, metaHeaders map[string][]string, authConfig *types.AuthConfig, outStream io.Writer) error { + start := time.Now() // Special case: "pull -a" may send an image name with a // trailing :. This is ugly, but let's not break API // compatibility. @@ -44,7 +46,9 @@ func (i *ImageService) PullImage(ctx context.Context, image, tag, os string, met } } - return i.pullImageWithReference(ctx, ref, os, metaHeaders, authConfig, outStream) + err = i.pullImageWithReference(ctx, ref, os, metaHeaders, authConfig, outStream) + imageActions.WithValues("pull").UpdateSince(start) + return err } func (i *ImageService) pullImageWithReference(ctx context.Context, ref reference.Named, os string, metaHeaders map[string][]string, authConfig *types.AuthConfig, outStream io.Writer) error { diff --git a/components/engine/daemon/images/image_push.go b/components/engine/daemon/images/image_push.go index a1a6d7f6d1..4c7be8d2e9 100644 --- a/components/engine/daemon/images/image_push.go +++ b/components/engine/daemon/images/image_push.go @@ -3,6 +3,7 @@ package images // import "github.com/docker/docker/daemon/images" import ( "context" "io" + "time" "github.com/docker/distribution/manifest/schema2" "github.com/docker/distribution/reference" @@ -14,6 +15,7 @@ import ( // PushImage initiates a push operation on the repository named localName. func (i *ImageService) PushImage(ctx context.Context, image, tag string, metaHeaders map[string][]string, authConfig *types.AuthConfig, outStream io.Writer) error { + start := time.Now() ref, err := reference.ParseNormalizedNamed(image) if err != nil { return err @@ -59,5 +61,6 @@ func (i *ImageService) PushImage(ctx context.Context, image, tag string, metaHea err = distribution.Push(ctx, ref, imagePushConfig) close(progressChan) <-writesDone + imageActions.WithValues("push").UpdateSince(start) return err } From 54e00c6a5cd20453b7bbbc1282049785ba1c92ae Mon Sep 17 00:00:00 2001 From: Brian Goff Date: Fri, 20 Apr 2018 10:48:54 -0400 Subject: [PATCH 10/13] Fix panic on daemon restart with running plugin Scenario: Daemon is ungracefully shutdown and leaves plugins running (no live-restore). Daemon comes back up. The next time a container tries to use that plugin it will cause a daemon panic because the plugin client is not set. This fixes that by ensuring that the plugin does get shutdown. Note, I do not think there would be any harm in just re-attaching to the running plugin instead of shutting it down, however historically we shut down plugins and containers when live-restore is not enabled. [kir@: consolidate code to deleteTaskAndContainer, a few minor nits] Signed-off-by: Brian Goff Signed-off-by: Kir Kolyshkin Upstream-commit: dbeb4329655e91dbe0e6574405937f03fabf3f2f Component: engine --- .../plugin/executor/containerd/containerd.go | 45 +++-- components/engine/plugin/manager.go | 22 ++- components/engine/plugin/manager_linux.go | 29 ++-- .../engine/plugin/manager_linux_test.go | 156 +++++++++++++++++- components/engine/plugin/manager_windows.go | 2 +- 5 files changed, 199 insertions(+), 55 deletions(-) diff --git a/components/engine/plugin/executor/containerd/containerd.go b/components/engine/plugin/executor/containerd/containerd.go index e490ef0a9e..e0267582ec 100644 --- a/components/engine/plugin/executor/containerd/containerd.go +++ b/components/engine/plugin/executor/containerd/containerd.go @@ -58,6 +58,19 @@ type Executor struct { exitHandler ExitHandler } +// deleteTaskAndContainer deletes plugin task and then plugin container from containerd +func deleteTaskAndContainer(ctx context.Context, cli Client, id string) { + _, _, err := cli.DeleteTask(ctx, id) + if err != nil && !errdefs.IsNotFound(err) { + logrus.WithError(err).WithField("id", id).Error("failed to delete plugin task from containerd") + } + + err = cli.Delete(ctx, id) + if err != nil && !errdefs.IsNotFound(err) { + logrus.WithError(err).WithField("id", id).Error("failed to delete plugin container from containerd") + } +} + // Create creates a new container func (e *Executor) Create(id string, spec specs.Spec, stdout, stderr io.WriteCloser) error { opts := runctypes.RuncOptions{ @@ -87,34 +100,21 @@ func (e *Executor) Create(id string, spec specs.Spec, stdout, stderr io.WriteClo _, err = e.client.Start(ctx, id, "", false, attachStreamsFunc(stdout, stderr)) if err != nil { - if _, _, err2 := e.client.DeleteTask(ctx, id); err2 != nil && !errdefs.IsNotFound(err2) { - logrus.WithError(err2).WithField("id", id).Warn("Received an error while attempting to clean up containerd plugin task after failed start") - } - if err2 := e.client.Delete(ctx, id); err2 != nil && !errdefs.IsNotFound(err2) { - logrus.WithError(err2).WithField("id", id).Warn("Received an error while attempting to clean up containerd plugin container after failed start") - } + deleteTaskAndContainer(ctx, e.client, id) } return err } // Restore restores a container -func (e *Executor) Restore(id string, stdout, stderr io.WriteCloser) error { +func (e *Executor) Restore(id string, stdout, stderr io.WriteCloser) (bool, error) { alive, _, err := e.client.Restore(context.Background(), id, attachStreamsFunc(stdout, stderr)) if err != nil && !errdefs.IsNotFound(err) { - return err + return false, err } if !alive { - _, _, err = e.client.DeleteTask(context.Background(), id) - if err != nil && !errdefs.IsNotFound(err) { - logrus.WithError(err).Errorf("failed to delete container plugin %s task from containerd", id) - } - - err = e.client.Delete(context.Background(), id) - if err != nil && !errdefs.IsNotFound(err) { - logrus.WithError(err).Errorf("failed to delete container plugin %s from containerd", id) - } + deleteTaskAndContainer(context.Background(), e.client, id) } - return nil + return alive, nil } // IsRunning returns if the container with the given id is running @@ -133,14 +133,7 @@ func (e *Executor) Signal(id string, signal int) error { func (e *Executor) ProcessEvent(id string, et libcontainerd.EventType, ei libcontainerd.EventInfo) error { switch et { case libcontainerd.EventExit: - // delete task and container - if _, _, err := e.client.DeleteTask(context.Background(), id); err != nil { - logrus.WithError(err).Errorf("failed to delete container plugin %s task from containerd", id) - } - - if err := e.client.Delete(context.Background(), id); err != nil { - logrus.WithError(err).Errorf("failed to delete container plugin %s from containerd", id) - } + deleteTaskAndContainer(context.Background(), e.client, id) return e.exitHandler.HandleExitEvent(ei.ContainerID) } return nil diff --git a/components/engine/plugin/manager.go b/components/engine/plugin/manager.go index 9c674f9545..c6f896129b 100644 --- a/components/engine/plugin/manager.go +++ b/components/engine/plugin/manager.go @@ -37,14 +37,14 @@ var validFullID = regexp.MustCompile(`^([a-f0-9]{64})$`) // Executor is the interface that the plugin manager uses to interact with for starting/stopping plugins type Executor interface { Create(id string, spec specs.Spec, stdout, stderr io.WriteCloser) error - Restore(id string, stdout, stderr io.WriteCloser) error IsRunning(id string) (bool, error) + Restore(id string, stdout, stderr io.WriteCloser) (alive bool, err error) Signal(id string, signal int) error } -func (pm *Manager) restorePlugin(p *v2.Plugin) error { +func (pm *Manager) restorePlugin(p *v2.Plugin, c *controller) error { if p.IsEnabled() { - return pm.restore(p) + return pm.restore(p, c) } return nil } @@ -143,12 +143,15 @@ func (pm *Manager) HandleExitEvent(id string) error { return err } - os.RemoveAll(filepath.Join(pm.config.ExecRoot, id)) + if err := os.RemoveAll(filepath.Join(pm.config.ExecRoot, id)); err != nil && !os.IsNotExist(err) { + logrus.WithError(err).WithField("id", id).Error("Could not remove plugin bundle dir") + } pm.mu.RLock() c := pm.cMap[p] if c.exitChan != nil { close(c.exitChan) + c.exitChan = nil // ignore duplicate events (containerd issue #2299) } restart := c.restart pm.mu.RUnlock() @@ -205,12 +208,15 @@ func (pm *Manager) reload() error { // todo: restore var wg sync.WaitGroup wg.Add(len(plugins)) for _, p := range plugins { - c := &controller{} // todo: remove this + c := &controller{exitChan: make(chan bool)} + pm.mu.Lock() pm.cMap[p] = c + pm.mu.Unlock() + go func(p *v2.Plugin) { defer wg.Done() - if err := pm.restorePlugin(p); err != nil { - logrus.Errorf("failed to restore plugin '%s': %s", p.Name(), err) + if err := pm.restorePlugin(p, c); err != nil { + logrus.WithError(err).WithField("id", p.GetID()).Error("Failed to restore plugin") return } @@ -248,7 +254,7 @@ func (pm *Manager) reload() error { // todo: restore if requiresManualRestore { // if liveRestore is not enabled, the plugin will be stopped now so we should enable it if err := pm.enable(p, c, true); err != nil { - logrus.Errorf("failed to enable plugin '%s': %s", p.Name(), err) + logrus.WithError(err).WithField("id", p.GetID()).Error("failed to enable plugin") } } }(p) diff --git a/components/engine/plugin/manager_linux.go b/components/engine/plugin/manager_linux.go index 0029ff7868..3c6f9c553a 100644 --- a/components/engine/plugin/manager_linux.go +++ b/components/engine/plugin/manager_linux.go @@ -79,7 +79,7 @@ func (pm *Manager) pluginPostStart(p *v2.Plugin, c *controller) error { client, err := plugins.NewClientWithTimeout(addr.Network()+"://"+addr.String(), nil, p.Timeout()) if err != nil { c.restart = false - shutdownPlugin(p, c, pm.executor) + shutdownPlugin(p, c.exitChan, pm.executor) return errors.WithStack(err) } @@ -106,7 +106,7 @@ func (pm *Manager) pluginPostStart(p *v2.Plugin, c *controller) error { c.restart = false // While restoring plugins, we need to explicitly set the state to disabled pm.config.Store.SetState(p, false) - shutdownPlugin(p, c, pm.executor) + shutdownPlugin(p, c.exitChan, pm.executor) return err } @@ -117,16 +117,15 @@ func (pm *Manager) pluginPostStart(p *v2.Plugin, c *controller) error { return pm.save(p) } -func (pm *Manager) restore(p *v2.Plugin) error { +func (pm *Manager) restore(p *v2.Plugin, c *controller) error { stdout, stderr := makeLoggerStreams(p.GetID()) - if err := pm.executor.Restore(p.GetID(), stdout, stderr); err != nil { + alive, err := pm.executor.Restore(p.GetID(), stdout, stderr) + if err != nil { return err } if pm.config.LiveRestoreEnabled { - c := &controller{} - if isRunning, _ := pm.executor.IsRunning(p.GetID()); !isRunning { - // plugin is not running, so follow normal startup procedure + if !alive { return pm.enable(p, c, true) } @@ -138,10 +137,16 @@ func (pm *Manager) restore(p *v2.Plugin) error { return pm.pluginPostStart(p, c) } + if alive { + // TODO(@cpuguy83): Should we always just re-attach to the running plugin instead of doing this? + c.restart = false + shutdownPlugin(p, c.exitChan, pm.executor) + } + return nil } -func shutdownPlugin(p *v2.Plugin, c *controller, executor Executor) { +func shutdownPlugin(p *v2.Plugin, ec chan bool, executor Executor) { pluginID := p.GetID() err := executor.Signal(pluginID, int(unix.SIGTERM)) @@ -149,7 +154,7 @@ func shutdownPlugin(p *v2.Plugin, c *controller, executor Executor) { logrus.Errorf("Sending SIGTERM to plugin failed with error: %v", err) } else { select { - case <-c.exitChan: + case <-ec: logrus.Debug("Clean shutdown of plugin") case <-time.After(time.Second * 10): logrus.Debug("Force shutdown plugin") @@ -157,7 +162,7 @@ func shutdownPlugin(p *v2.Plugin, c *controller, executor Executor) { logrus.Errorf("Sending SIGKILL to plugin failed with error: %v", err) } select { - case <-c.exitChan: + case <-ec: logrus.Debug("SIGKILL plugin shutdown") case <-time.After(time.Second * 10): logrus.Debug("Force shutdown plugin FAILED") @@ -172,7 +177,7 @@ func (pm *Manager) disable(p *v2.Plugin, c *controller) error { } c.restart = false - shutdownPlugin(p, c, pm.executor) + shutdownPlugin(p, c.exitChan, pm.executor) pm.config.Store.SetState(p, false) return pm.save(p) } @@ -191,7 +196,7 @@ func (pm *Manager) Shutdown() { } if pm.executor != nil && p.IsEnabled() { c.restart = false - shutdownPlugin(p, c, pm.executor) + shutdownPlugin(p, c.exitChan, pm.executor) } } if err := mount.RecursiveUnmount(pm.config.Root); err != nil { diff --git a/components/engine/plugin/manager_linux_test.go b/components/engine/plugin/manager_linux_test.go index d4199c80da..740efd7a3a 100644 --- a/components/engine/plugin/manager_linux_test.go +++ b/components/engine/plugin/manager_linux_test.go @@ -3,12 +3,14 @@ package plugin // import "github.com/docker/docker/plugin" import ( "io" "io/ioutil" + "net" "os" "path/filepath" "testing" "github.com/docker/docker/api/types" "github.com/docker/docker/pkg/mount" + "github.com/docker/docker/pkg/stringid" "github.com/docker/docker/pkg/system" "github.com/docker/docker/plugin/v2" "github.com/gotestyourself/gotestyourself/skip" @@ -59,7 +61,7 @@ func TestManagerWithPluginMounts(t *testing.T) { t.Fatal(err) } - if err := m.Remove(p1.Name(), &types.PluginRmConfig{ForceRemove: true}); err != nil { + if err := m.Remove(p1.GetID(), &types.PluginRmConfig{ForceRemove: true}); err != nil { t.Fatal(err) } if mounted, err := mount.Mounted(p2Mount); !mounted || err != nil { @@ -68,17 +70,18 @@ func TestManagerWithPluginMounts(t *testing.T) { } func newTestPlugin(t *testing.T, name, cap, root string) *v2.Plugin { - rootfs := filepath.Join(root, name) + id := stringid.GenerateNonCryptoID() + rootfs := filepath.Join(root, id) if err := os.MkdirAll(rootfs, 0755); err != nil { t.Fatal(err) } - p := v2.Plugin{PluginObj: types.Plugin{Name: name}} + p := v2.Plugin{PluginObj: types.Plugin{ID: id, Name: name}} p.Rootfs = rootfs iType := types.PluginInterfaceType{Capability: cap, Prefix: "docker", Version: "1.0"} - i := types.PluginConfigInterface{Socket: "plugins.sock", Types: []types.PluginInterfaceType{iType}} + i := types.PluginConfigInterface{Socket: "plugin.sock", Types: []types.PluginInterfaceType{iType}} p.PluginObj.Config.Interface = i - p.PluginObj.ID = name + p.PluginObj.ID = id return &p } @@ -90,8 +93,8 @@ func (e *simpleExecutor) Create(id string, spec specs.Spec, stdout, stderr io.Wr return errors.New("Create failed") } -func (e *simpleExecutor) Restore(id string, stdout, stderr io.WriteCloser) error { - return nil +func (e *simpleExecutor) Restore(id string, stdout, stderr io.WriteCloser) (bool, error) { + return false, nil } func (e *simpleExecutor) IsRunning(id string) (bool, error) { @@ -133,7 +136,144 @@ func TestCreateFailed(t *testing.T) { t.Fatalf("expected Create failed error, got %v", err) } - if err := m.Remove(p.Name(), &types.PluginRmConfig{ForceRemove: true}); err != nil { + if err := m.Remove(p.GetID(), &types.PluginRmConfig{ForceRemove: true}); err != nil { t.Fatal(err) } } + +type executorWithRunning struct { + m *Manager + root string + exitChans map[string]chan struct{} +} + +func (e *executorWithRunning) Create(id string, spec specs.Spec, stdout, stderr io.WriteCloser) error { + sockAddr := filepath.Join(e.root, id, "plugin.sock") + ch := make(chan struct{}) + if e.exitChans == nil { + e.exitChans = make(map[string]chan struct{}) + } + e.exitChans[id] = ch + listenTestPlugin(sockAddr, ch) + return nil +} + +func (e *executorWithRunning) IsRunning(id string) (bool, error) { + return true, nil +} +func (e *executorWithRunning) Restore(id string, stdout, stderr io.WriteCloser) (bool, error) { + return true, nil +} + +func (e *executorWithRunning) Signal(id string, signal int) error { + ch := e.exitChans[id] + ch <- struct{}{} + <-ch + e.m.HandleExitEvent(id) + return nil +} + +func TestPluginAlreadyRunningOnStartup(t *testing.T) { + t.Parallel() + + root, err := ioutil.TempDir("", t.Name()) + if err != nil { + t.Fatal(err) + } + defer system.EnsureRemoveAll(root) + + for _, test := range []struct { + desc string + config ManagerConfig + }{ + { + desc: "live-restore-disabled", + config: ManagerConfig{ + LogPluginEvent: func(_, _, _ string) {}, + }, + }, + { + desc: "live-restore-enabled", + config: ManagerConfig{ + LogPluginEvent: func(_, _, _ string) {}, + LiveRestoreEnabled: true, + }, + }, + } { + t.Run(test.desc, func(t *testing.T) { + config := test.config + desc := test.desc + t.Parallel() + + p := newTestPlugin(t, desc, desc, config.Root) + p.PluginObj.Enabled = true + + // Need a short-ish path here so we don't run into unix socket path length issues. + config.ExecRoot, err = ioutil.TempDir("", "plugintest") + + executor := &executorWithRunning{root: config.ExecRoot} + config.CreateExecutor = func(m *Manager) (Executor, error) { executor.m = m; return executor, nil } + + if err := executor.Create(p.GetID(), specs.Spec{}, nil, nil); err != nil { + t.Fatal(err) + } + + root := filepath.Join(root, desc) + config.Root = filepath.Join(root, "manager") + if err := os.MkdirAll(filepath.Join(config.Root, p.GetID()), 0755); err != nil { + t.Fatal(err) + } + + if !p.IsEnabled() { + t.Fatal("plugin should be enabled") + } + if err := (&Manager{config: config}).save(p); err != nil { + t.Fatal(err) + } + + s := NewStore() + config.Store = s + if err != nil { + t.Fatal(err) + } + defer system.EnsureRemoveAll(config.ExecRoot) + + m, err := NewManager(config) + if err != nil { + t.Fatal(err) + } + defer m.Shutdown() + + p = s.GetAll()[p.GetID()] // refresh `p` with what the manager knows + if p.Client() == nil { + t.Fatal("plugin client should not be nil") + } + }) + } +} + +func listenTestPlugin(sockAddr string, exit chan struct{}) (net.Listener, error) { + if err := os.MkdirAll(filepath.Dir(sockAddr), 0755); err != nil { + return nil, err + } + l, err := net.Listen("unix", sockAddr) + if err != nil { + return nil, err + } + go func() { + for { + conn, err := l.Accept() + if err != nil { + return + } + conn.Close() + } + }() + go func() { + <-exit + l.Close() + os.Remove(sockAddr) + exit <- struct{}{} + }() + return l, nil +} diff --git a/components/engine/plugin/manager_windows.go b/components/engine/plugin/manager_windows.go index 9fafea5c22..90cc52c992 100644 --- a/components/engine/plugin/manager_windows.go +++ b/components/engine/plugin/manager_windows.go @@ -19,7 +19,7 @@ func (pm *Manager) disable(p *v2.Plugin, c *controller) error { return fmt.Errorf("Not implemented") } -func (pm *Manager) restore(p *v2.Plugin) error { +func (pm *Manager) restore(p *v2.Plugin, c *controller) error { return fmt.Errorf("Not implemented") } From 2ca4487c2680af18a6c224d3e3ccf04e3688bdae Mon Sep 17 00:00:00 2001 From: fanjiyun Date: Tue, 30 Jan 2018 16:02:59 +0800 Subject: [PATCH 11/13] When id is empty for overlay2/overlay, do not remove the directories. Signed-off-by: fanjiyun Signed-off-by: Sebastiaan van Stijn Upstream-commit: 0e8f96e31724a7bb49d0ade9acec116f68c85c74 Component: engine --- components/engine/daemon/graphdriver/overlay/overlay.go | 3 +++ components/engine/daemon/graphdriver/overlay2/overlay.go | 7 ++++++- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/components/engine/daemon/graphdriver/overlay/overlay.go b/components/engine/daemon/graphdriver/overlay/overlay.go index 62dc26a871..0c2167f083 100644 --- a/components/engine/daemon/graphdriver/overlay/overlay.go +++ b/components/engine/daemon/graphdriver/overlay/overlay.go @@ -366,6 +366,9 @@ func (d *Driver) dir(id string) string { // Remove cleans the directories that are created for this id. func (d *Driver) Remove(id string) error { + if id == "" { + return fmt.Errorf("refusing to remove the directories: id is empty") + } d.locker.Lock(id) defer d.locker.Unlock(id) return system.EnsureRemoveAll(d.dir(id)) diff --git a/components/engine/daemon/graphdriver/overlay2/overlay.go b/components/engine/daemon/graphdriver/overlay2/overlay.go index 2b52adb858..5108a2c055 100644 --- a/components/engine/daemon/graphdriver/overlay2/overlay.go +++ b/components/engine/daemon/graphdriver/overlay2/overlay.go @@ -513,12 +513,17 @@ func (d *Driver) getLowerDirs(id string) ([]string, error) { // Remove cleans the directories that are created for this id. func (d *Driver) Remove(id string) error { + if id == "" { + return fmt.Errorf("refusing to remove the directories: id is empty") + } d.locker.Lock(id) defer d.locker.Unlock(id) dir := d.dir(id) lid, err := ioutil.ReadFile(path.Join(dir, "link")) if err == nil { - if err := os.RemoveAll(path.Join(d.home, linkDir, string(lid))); err != nil { + if len(lid) == 0 { + logrus.WithField("storage-driver", "overlay2").Errorf("refusing to remove empty link for layer %v", id) + } else if err := os.RemoveAll(path.Join(d.home, linkDir, string(lid))); err != nil { logrus.WithField("storage-driver", "overlay2").Debugf("Failed to remove link: %v", err) } } From 3f03982f9a321d1d357bc186d1f99cb4e1a80854 Mon Sep 17 00:00:00 2001 From: Kunal Tyagi Date: Fri, 8 Jun 2018 10:10:09 +0900 Subject: [PATCH 12/13] Allow vim be case insensitive for D in dockerfile Signed-off-by: Kunal Tyagi Upstream-commit: 6b8dab2181097c83670b65c2f76f83307058c987 Component: engine --- components/engine/contrib/syntax/vim/ftdetect/dockerfile.vim | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/engine/contrib/syntax/vim/ftdetect/dockerfile.vim b/components/engine/contrib/syntax/vim/ftdetect/dockerfile.vim index f7a962e073..a21dd14095 100644 --- a/components/engine/contrib/syntax/vim/ftdetect/dockerfile.vim +++ b/components/engine/contrib/syntax/vim/ftdetect/dockerfile.vim @@ -1 +1 @@ -au BufNewFile,BufRead [Dd]ockerfile,Dockerfile.*,*.Dockerfile set filetype=dockerfile +au BufNewFile,BufRead [Dd]ockerfile,[Dd]ockerfile.*,*.[Dd]ockerfile set filetype=dockerfile From 1fd604f1f8451757d65d40c09654ef7bc9a8ed39 Mon Sep 17 00:00:00 2001 From: Francesco Mari Date: Fri, 8 Jun 2018 16:09:46 +0200 Subject: [PATCH 13/13] Fix link to Docker Toolbox Signed-off-by: Francesco Mari Upstream-commit: a045b027bf8f9e19fd53c2b33399fad901b1dd34 Component: engine --- components/engine/docs/contributing/set-up-dev-env.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/engine/docs/contributing/set-up-dev-env.md b/components/engine/docs/contributing/set-up-dev-env.md index 311edf8951..3d56c0b8c7 100644 --- a/components/engine/docs/contributing/set-up-dev-env.md +++ b/components/engine/docs/contributing/set-up-dev-env.md @@ -93,7 +93,7 @@ can take over 15 minutes to complete. 1. Open a terminal. - For [Docker Toolbox](../../toolbox/overview.md) users, use `docker-machine status your_vm_name` to make sure your VM is running. You + For [Docker Toolbox](https://github.com/docker/toolbox) users, use `docker-machine status your_vm_name` to make sure your VM is running. You may need to run `eval "$(docker-machine env your_vm_name)"` to initialize your shell environment. If you use Docker for Mac or Docker for Windows, you do not need to use Docker Machine.