diff --git a/components/engine/vendor.conf b/components/engine/vendor.conf index 1f2ffc1f29..22004d2916 100644 --- a/components/engine/vendor.conf +++ b/components/engine/vendor.conf @@ -110,7 +110,7 @@ github.com/tonistiigi/fifo 1405643975692217d6720f8b54aeee1bf2cd5cf4 github.com/stevvooe/continuity cd7a8e21e2b6f84799f5dd4b65faf49c8d3ee02d # cluster -github.com/docker/swarmkit 941a01844b89c56aa61086fecb167ab3af1de22b +github.com/docker/swarmkit 872861d2ae46958af7ead1d5fffb092c73afbaf0 github.com/gogo/protobuf v0.4 github.com/cloudflare/cfssl 7fb22c8cba7ecaf98e4082d22d65800cf45e042a github.com/google/certificate-transparency d90e65c3a07988180c5b1ece71791c0b6506826e diff --git a/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/networkallocator.go b/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/networkallocator.go index ab1b6d6a47..53f9ffbeee 100644 --- a/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/networkallocator.go +++ b/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/networkallocator.go @@ -574,6 +574,7 @@ func (na *cnmNetworkAllocator) releaseEndpoints(networks []*api.NetworkAttachmen // allocate virtual IP for a single endpoint attachment of the service. func (na *cnmNetworkAllocator) allocateVIP(vip *api.Endpoint_VirtualIP) error { + var opts map[string]string localNet := na.getNetwork(vip.NetworkID) if localNet == nil { return errors.New("networkallocator: could not find local network state") @@ -603,9 +604,13 @@ func (na *cnmNetworkAllocator) allocateVIP(vip *api.Endpoint_VirtualIP) error { return err } } + if localNet.nw.IPAM != nil && localNet.nw.IPAM.Driver != nil { + // set ipam allocation method to serial + opts = setIPAMSerialAlloc(localNet.nw.IPAM.Driver.Options) + } for _, poolID := range localNet.pools { - ip, _, err := ipam.RequestAddress(poolID, addr, nil) + ip, _, err := ipam.RequestAddress(poolID, addr, opts) if err != nil && err != ipamapi.ErrNoAvailableIPs && err != ipamapi.ErrIPOutOfRange { return errors.Wrap(err, "could not allocate VIP from IPAM") } @@ -657,6 +662,7 @@ func (na *cnmNetworkAllocator) deallocateVIP(vip *api.Endpoint_VirtualIP) error // allocate the IP addresses for a single network attachment of the task. func (na *cnmNetworkAllocator) allocateNetworkIPs(nAttach *api.NetworkAttachment) error { var ip *net.IPNet + var opts map[string]string ipam, _, _, err := na.resolveIPAM(nAttach.Network) if err != nil { @@ -686,11 +692,16 @@ func (na *cnmNetworkAllocator) allocateNetworkIPs(nAttach *api.NetworkAttachment } } } + // Set the ipam options if the network has an ipam driver. + if localNet.nw.IPAM != nil && localNet.nw.IPAM.Driver != nil { + // set ipam allocation method to serial + opts = setIPAMSerialAlloc(localNet.nw.IPAM.Driver.Options) + } for _, poolID := range localNet.pools { var err error - ip, _, err = ipam.RequestAddress(poolID, addr, nil) + ip, _, err = ipam.RequestAddress(poolID, addr, opts) if err != nil && err != ipamapi.ErrNoAvailableIPs && err != ipamapi.ErrIPOutOfRange { return errors.Wrap(err, "could not allocate IP from IPAM") } @@ -918,8 +929,16 @@ func (na *cnmNetworkAllocator) allocatePools(n *api.Network) (map[string]string, } gwIP.IP = ip } + if dOptions == nil { + dOptions = make(map[string]string) + } + dOptions[ipamapi.RequestAddressType] = netlabel.Gateway + // set ipam allocation method to serial + dOptions = setIPAMSerialAlloc(dOptions) + defer delete(dOptions, ipamapi.RequestAddressType) + if ic.Gateway != "" || gwIP == nil { - gwIP, _, err = ipam.RequestAddress(poolID, net.ParseIP(ic.Gateway), map[string]string{ipamapi.RequestAddressType: netlabel.Gateway}) + gwIP, _, err = ipam.RequestAddress(poolID, net.ParseIP(ic.Gateway), dOptions) if err != nil { // Rollback by releasing all the resources allocated so far. releasePools(ipam, ipamConfigs[:i], pools) @@ -980,3 +999,14 @@ func IsBuiltInDriver(name string) bool { } return false } + +// setIPAMSerialAlloc sets the ipam allocation method to serial +func setIPAMSerialAlloc(opts map[string]string) map[string]string { + if opts == nil { + opts = make(map[string]string) + } + if _, ok := opts[ipamapi.AllocSerialPrefix]; !ok { + opts[ipamapi.AllocSerialPrefix] = "true" + } + return opts +} diff --git a/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/portallocator.go b/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/portallocator.go index b09ac47c79..19dcbec772 100644 --- a/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/portallocator.go +++ b/components/engine/vendor/github.com/docker/swarmkit/manager/allocator/cnmallocator/portallocator.go @@ -382,7 +382,7 @@ func (ps *portSpace) allocate(p *api.PortConfig) (err error) { } // Check out an arbitrary port from dynamic port space. - swarmPort, err := ps.dynamicPortSpace.GetID() + swarmPort, err := ps.dynamicPortSpace.GetID(true) if err != nil { return } diff --git a/components/engine/vendor/github.com/docker/swarmkit/manager/state/raft/raft.go b/components/engine/vendor/github.com/docker/swarmkit/manager/state/raft/raft.go index ec2c9c0f5e..afdf2ca4eb 100644 --- a/components/engine/vendor/github.com/docker/swarmkit/manager/state/raft/raft.go +++ b/components/engine/vendor/github.com/docker/swarmkit/manager/state/raft/raft.go @@ -542,6 +542,7 @@ func (n *Node) Run(ctx context.Context) error { n.done() }() + // Flag that indicates if this manager node is *currently* the raft leader. wasLeader := false transferLeadershipLimit := rate.NewLimiter(rate.Every(time.Minute), 1) @@ -563,10 +564,13 @@ func (n *Node) Run(ctx context.Context) error { return errors.Wrap(err, "failed to save entries to storage") } + // If the memory store lock has been held for too long, + // transferring leadership is an easy way to break out of it. if wasLeader && (rd.SoftState == nil || rd.SoftState.RaftState == raft.StateLeader) && n.memoryStore.Wedged() && transferLeadershipLimit.Allow() { + log.G(ctx).Error("Attempting to transfer leadership") if !n.opts.DisableStackDump { signal.DumpStacks("") } @@ -612,6 +616,8 @@ func (n *Node) Run(ctx context.Context) error { if rd.SoftState != nil { if wasLeader && rd.SoftState.RaftState != raft.StateLeader { wasLeader = false + log.G(ctx).Error("soft state changed, node no longer a leader, resetting and cancelling all waits") + if atomic.LoadUint32(&n.signalledLeadership) == 1 { atomic.StoreUint32(&n.signalledLeadership, 0) n.leadershipBroadcast.Publish(IsFollower) @@ -630,6 +636,7 @@ func (n *Node) Run(ctx context.Context) error { // cancelAll, or by its own check of signalledLeadership. n.wait.cancelAll() } else if !wasLeader && rd.SoftState.RaftState == raft.StateLeader { + // Node just became a leader. wasLeader = true } } @@ -1478,7 +1485,7 @@ func (n *Node) registerNode(node *api.RaftMember) error { return nil } -// ProposeValue calls Propose on the raft and waits +// ProposeValue calls Propose on the underlying raft library(etcd/raft) and waits // on the commit log action before returning a result func (n *Node) ProposeValue(ctx context.Context, storeAction []api.StoreAction, cb func()) error { ctx, cancel := n.WithContext(ctx) @@ -1654,11 +1661,14 @@ func (n *Node) saveToStorage( return nil } -// processInternalRaftRequest sends a message to nodes participating -// in the raft to apply a log entry and then waits for it to be applied -// on the server. It will block until the update is performed, there is -// an error or until the raft node finalizes all the proposals on node -// shutdown. +// processInternalRaftRequest proposes a value to be appended to the raft log. +// It calls Propose() on etcd/raft, which calls back into the raft FSM, +// which then sends a message to each of the participating nodes +// in the raft group to apply a log entry and then waits for it to be applied +// on this node. It will block until the this node: +// 1. Gets the necessary replies back from the participating nodes and also performs the commit itself, or +// 2. There is an error, or +// 3. Until the raft node finalizes all the proposals on node shutdown. func (n *Node) processInternalRaftRequest(ctx context.Context, r *api.InternalRaftRequest, cb func()) (proto.Message, error) { n.stopMu.RLock() if !n.IsMember() { @@ -1679,6 +1689,7 @@ func (n *Node) processInternalRaftRequest(ctx context.Context, r *api.InternalRa // Do this check after calling register to avoid a race. if atomic.LoadUint32(&n.signalledLeadership) != 1 { + log.G(ctx).Error("node is no longer leader, aborting propose") n.wait.cancel(r.ID) return nil, ErrLostLeadership } @@ -1703,14 +1714,23 @@ func (n *Node) processInternalRaftRequest(ctx context.Context, r *api.InternalRa select { case x, ok := <-ch: if !ok { + // Wait notification channel was closed. This should only happen if the wait was cancelled. + log.G(ctx).Error("wait cancelled") + if atomic.LoadUint32(&n.signalledLeadership) == 1 { + log.G(ctx).Error("wait cancelled but node is still a leader") + } return nil, ErrLostLeadership } return x.(proto.Message), nil case <-waitCtx.Done(): n.wait.cancel(r.ID) - // if channel is closed, wait item was canceled, otherwise it was triggered + // If we can read from the channel, wait item was triggered. Otherwise it was cancelled. x, ok := <-ch if !ok { + log.G(ctx).WithError(waitCtx.Err()).Error("wait context cancelled") + if atomic.LoadUint32(&n.signalledLeadership) == 1 { + log.G(ctx).Error("wait context cancelled but node is still a leader") + } return nil, ErrLostLeadership } return x.(proto.Message), nil @@ -1779,21 +1799,26 @@ func (n *Node) processEntry(ctx context.Context, entry raftpb.Entry) error { } if !n.wait.trigger(r.ID, r) { + log.G(ctx).Errorf("wait not found for raft request id %x", r.ID) + // There was no wait on this ID, meaning we don't have a // transaction in progress that would be committed to the // memory store by the "trigger" call. Either a different node // wrote this to raft, or we wrote it before losing the leader - // position and cancelling the transaction. Create a new - // transaction to commit the data. + // position and cancelling the transaction. This entry still needs + // to be committed since other nodes have already committed it. + // Create a new transaction to commit this entry. // It should not be possible for processInternalRaftRequest // to be running in this situation, but out of caution we // cancel any current invocations to avoid a deadlock. + // TODO(anshul) This call is likely redundant, remove after consideration. n.wait.cancelAll() err := n.memoryStore.ApplyStoreActions(r.Action) if err != nil { log.G(ctx).WithError(err).Error("failed to apply actions from raft") + // TODO(anshul) return err here ? } } return nil diff --git a/components/engine/vendor/github.com/docker/swarmkit/manager/state/store/memory.go b/components/engine/vendor/github.com/docker/swarmkit/manager/state/store/memory.go index 62ab9279d4..01245a6966 100644 --- a/components/engine/vendor/github.com/docker/swarmkit/manager/state/store/memory.go +++ b/components/engine/vendor/github.com/docker/swarmkit/manager/state/store/memory.go @@ -83,8 +83,7 @@ func register(os ObjectStoreConfig) { schema.Tables[os.Table.Name] = os.Table } -// timedMutex wraps a sync.Mutex, and keeps track of how long it has been -// locked. +// timedMutex wraps a sync.Mutex, and keeps track of when it was locked. type timedMutex struct { sync.Mutex lockedAt atomic.Value diff --git a/components/engine/vendor/github.com/docker/swarmkit/vendor.conf b/components/engine/vendor/github.com/docker/swarmkit/vendor.conf index abb1a6c247..8949ea01fc 100644 --- a/components/engine/vendor/github.com/docker/swarmkit/vendor.conf +++ b/components/engine/vendor/github.com/docker/swarmkit/vendor.conf @@ -24,7 +24,7 @@ github.com/docker/go-connections 3ede32e2033de7505e6500d6c868c2b9ed9f169d github.com/docker/go-events 9461782956ad83b30282bf90e31fa6a70c255ba9 github.com/docker/go-units 954fed01cc617c55d838fa2230073f2cb17386c8 github.com/docker/libkv 9fd56606e928ff1f309808f5d5a0b7a2ef73f9a8 -github.com/docker/libnetwork 19ac3ea7f52bb46e0eb10669756cdae0c441a5b1 +github.com/docker/libnetwork 21544598c53fa36a3c771a8725c643dd2340f845 github.com/docker/libtrust 9cbd2a1374f46905c68a4eb3694a130610adc62a github.com/opencontainers/runc d40db12e72a40109dfcf28539f5ee0930d2f0277 github.com/opencontainers/go-digest 21dfd564fd89c944783d00d069f33e3e7123c448