limiter.go 4.12 KB
Newer Older
Jeromy's avatar
Jeromy committed
1 2 3
package swarm

import (
4
	"context"
Jeromy's avatar
Jeromy committed
5 6
	"sync"

Jeromy's avatar
Jeromy committed
7 8
	addrutil "github.com/libp2p/go-addr-util"
	peer "github.com/libp2p/go-libp2p-peer"
Steven Allen's avatar
Steven Allen committed
9
	transport "github.com/libp2p/go-libp2p-transport"
Jeromy's avatar
Jeromy committed
10
	ma "github.com/multiformats/go-multiaddr"
Jeromy's avatar
Jeromy committed
11 12 13
)

type dialResult struct {
Steven Allen's avatar
Steven Allen committed
14
	Conn transport.Conn
Jeromy's avatar
Jeromy committed
15
	Addr ma.Multiaddr
Jeromy's avatar
Jeromy committed
16 17 18 19 20 21 22 23 24 25 26
	Err  error
}

type dialJob struct {
	addr    ma.Multiaddr
	peer    peer.ID
	ctx     context.Context
	resp    chan dialResult
	success bool
}

Jeromy's avatar
Jeromy committed
27 28 29 30 31 32 33 34 35
func (dj *dialJob) cancelled() bool {
	select {
	case <-dj.ctx.Done():
		return true
	default:
		return false
	}
}

Jeromy's avatar
Jeromy committed
36
type dialLimiter struct {
Łukasz Magiera's avatar
Łukasz Magiera committed
37 38
	lk sync.Mutex

Jeromy's avatar
Jeromy committed
39 40 41 42
	fdConsuming int
	fdLimit     int
	waitingOnFd []*dialJob

Łukasz Magiera's avatar
Łukasz Magiera committed
43
	dialFunc dialfunc
Jeromy's avatar
Jeromy committed
44 45 46 47 48 49

	activePerPeer      map[peer.ID]int
	perPeerLimit       int
	waitingOnPeerLimit map[peer.ID][]*dialJob
}

Steven Allen's avatar
Steven Allen committed
50
type dialfunc func(context.Context, peer.ID, ma.Multiaddr) (transport.Conn, error)
Jeromy's avatar
Jeromy committed
51 52

func newDialLimiter(df dialfunc) *dialLimiter {
Steven Allen's avatar
Steven Allen committed
53
	return newDialLimiterWithParams(df, ConcurrentFdDials, DefaultPerPeerRateLimit)
Jeromy's avatar
Jeromy committed
54 55
}

56
func newDialLimiterWithParams(df dialfunc, fdLimit, perPeerLimit int) *dialLimiter {
Jeromy's avatar
Jeromy committed
57
	return &dialLimiter{
58 59
		fdLimit:            fdLimit,
		perPeerLimit:       perPeerLimit,
Jeromy's avatar
Jeromy committed
60 61 62 63 64 65
		waitingOnPeerLimit: make(map[peer.ID][]*dialJob),
		activePerPeer:      make(map[peer.ID]int),
		dialFunc:           df,
	}
}

Łukasz Magiera's avatar
Łukasz Magiera committed
66 67 68 69 70 71 72 73 74 75 76
// freeFDToken frees FD token and if there are any schedules another waiting dialJob
// in it's place
func (dl *dialLimiter) freeFDToken() {
	dl.fdConsuming--

	if len(dl.waitingOnFd) > 0 {
		next := dl.waitingOnFd[0]
		dl.waitingOnFd[0] = nil // clear out memory
		dl.waitingOnFd = dl.waitingOnFd[1:]
		if len(dl.waitingOnFd) == 0 {
			dl.waitingOnFd = nil // clear out memory
Jeromy's avatar
Jeromy committed
77
		}
Łukasz Magiera's avatar
Łukasz Magiera committed
78 79 80 81
		dl.fdConsuming++

		// we already have activePerPeer token at this point so we can just dial
		go dl.executeDial(next)
Jeromy's avatar
Jeromy committed
82
	}
Łukasz Magiera's avatar
Łukasz Magiera committed
83
}
Jeromy's avatar
Jeromy committed
84

Łukasz Magiera's avatar
Łukasz Magiera committed
85
func (dl *dialLimiter) freePeerToken(dj *dialJob) {
Jeromy's avatar
Jeromy committed
86 87 88 89 90 91 92 93 94 95
	// release tokens in reverse order than we take them
	dl.activePerPeer[dj.peer]--
	if dl.activePerPeer[dj.peer] == 0 {
		delete(dl.activePerPeer, dj.peer)
	}

	waitlist := dl.waitingOnPeerLimit[dj.peer]
	if !dj.success && len(waitlist) > 0 {
		next := waitlist[0]
		if len(waitlist) == 1 {
Łukasz Magiera's avatar
Łukasz Magiera committed
96
			delete(dl.waitingOnPeerLimit, next.peer)
Jeromy's avatar
Jeromy committed
97
		} else {
98
			waitlist[0] = nil // clear out memory
Łukasz Magiera's avatar
Łukasz Magiera committed
99
			dl.waitingOnPeerLimit[next.peer] = waitlist[1:]
Jeromy's avatar
Jeromy committed
100 101
		}

Łukasz Magiera's avatar
Łukasz Magiera committed
102
		dl.activePerPeer[next.peer]++ // just kidding, we still want this token
103

Łukasz Magiera's avatar
Łukasz Magiera committed
104
		dl.addCheckFdLimit(next)
Jeromy's avatar
Jeromy committed
105 106 107
	}
}

Łukasz Magiera's avatar
Łukasz Magiera committed
108 109 110
func (dl *dialLimiter) finishedDial(dj *dialJob) {
	dl.lk.Lock()
	defer dl.lk.Unlock()
Jeromy's avatar
Jeromy committed
111

Łukasz Magiera's avatar
Łukasz Magiera committed
112 113
	if addrutil.IsFDCostlyTransport(dj.addr) {
		dl.freeFDToken()
Jeromy's avatar
Jeromy committed
114 115
	}

Łukasz Magiera's avatar
Łukasz Magiera committed
116 117 118 119
	dl.freePeerToken(dj)
}

func (dl *dialLimiter) addCheckFdLimit(dj *dialJob) {
Jeromy's avatar
Jeromy committed
120 121 122 123 124 125 126 127 128 129
	if addrutil.IsFDCostlyTransport(dj.addr) {
		if dl.fdConsuming >= dl.fdLimit {
			dl.waitingOnFd = append(dl.waitingOnFd, dj)
			return
		}

		// take token
		dl.fdConsuming++
	}

Jeromy's avatar
Jeromy committed
130 131 132
	go dl.executeDial(dj)
}

Łukasz Magiera's avatar
Łukasz Magiera committed
133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153
func (dl *dialLimiter) addCheckPeerLimit(dj *dialJob) {
	if dl.activePerPeer[dj.peer] >= dl.perPeerLimit {
		wlist := dl.waitingOnPeerLimit[dj.peer]
		dl.waitingOnPeerLimit[dj.peer] = append(wlist, dj)
		return
	}
	dl.activePerPeer[dj.peer]++

	dl.addCheckFdLimit(dj)
}

// AddDialJob tries to take the needed tokens for starting the given dial job.
// If it acquires all needed tokens, it immediately starts the dial, otherwise
// it will put it on the waitlist for the requested token.
func (dl *dialLimiter) AddDialJob(dj *dialJob) {
	dl.lk.Lock()
	defer dl.lk.Unlock()

	dl.addCheckPeerLimit(dj)
}

154
func (dl *dialLimiter) clearAllPeerDials(p peer.ID) {
Łukasz Magiera's avatar
Łukasz Magiera committed
155 156
	dl.lk.Lock()
	defer dl.lk.Unlock()
157
	delete(dl.waitingOnPeerLimit, p)
Łukasz Magiera's avatar
Łukasz Magiera committed
158
	// NB: the waitingOnFd list doesn't need to be cleaned out here, we will
159 160 161 162
	// remove them as we encounter them because they are 'cancelled' at this
	// point
}

Jeromy's avatar
Jeromy committed
163 164 165 166 167
// executeDial calls the dialFunc, and reports the result through the response
// channel when finished. Once the response is sent it also releases all tokens
// it held during the dial.
func (dl *dialLimiter) executeDial(j *dialJob) {
	defer dl.finishedDial(j)
Jeromy's avatar
Jeromy committed
168 169 170 171
	if j.cancelled() {
		return
	}

Jeromy's avatar
Jeromy committed
172 173
	con, err := dl.dialFunc(j.ctx, j.peer, j.addr)
	select {
Jeromy's avatar
Jeromy committed
174
	case j.resp <- dialResult{Conn: con, Addr: j.addr, Err: err}:
Jeromy's avatar
Jeromy committed
175
	case <-j.ctx.Done():
176 177 178
		if err == nil {
			con.Close()
		}
Jeromy's avatar
Jeromy committed
179 180
	}
}