re-vndr swarmkit to 872861d
Signed-off-by: Andrew Hsu <andrewhsu@docker.com>
This commit is contained in:
@@ -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
|
||||
|
||||
Generated
Vendored
+33
-3
@@ -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
|
||||
}
|
||||
|
||||
Generated
Vendored
+1
-1
@@ -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
|
||||
}
|
||||
|
||||
+34
-9
@@ -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
|
||||
|
||||
Generated
Vendored
+1
-2
@@ -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
|
||||
|
||||
+1
-1
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user