-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathfront.go
More file actions
272 lines (240 loc) · 6.52 KB
/
Copy pathfront.go
File metadata and controls
272 lines (240 loc) · 6.52 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
package domainfront
import (
"cmp"
"context"
"math/rand/v2"
"slices"
"sync"
"time"
)
// front is a single domain front candidate with its runtime state.
type front struct {
Masquerade
LastSucceeded time.Time `json:"LastSucceeded"`
ProviderID string `json:"ProviderID"`
mu sync.RWMutex
}
func newFront(m *Masquerade, providerID string) *front {
return &front{
Masquerade: *m,
ProviderID: providerID,
}
}
func (f *front) lastSucceededTime() time.Time {
f.mu.RLock()
defer f.mu.RUnlock()
return f.LastSucceeded
}
func (f *front) setLastSucceeded(t time.Time) {
f.mu.Lock()
defer f.mu.Unlock()
f.LastSucceeded = t
}
func (f *front) markSucceeded() {
f.mu.Lock()
defer f.mu.Unlock()
f.LastSucceeded = time.Now()
}
func (f *front) markFailed() {
f.mu.Lock()
defer f.mu.Unlock()
f.LastSucceeded = time.Time{}
}
func (f *front) isSucceeding() bool {
f.mu.RLock()
defer f.mu.RUnlock()
return !f.LastSucceeded.IsZero()
}
// frontKey identifies a front by its provider, domain, and IP.
type frontKey struct{ provider, domain, ip string }
func (f *front) key() frontKey {
return frontKey{f.ProviderID, f.Domain, f.IpAddress}
}
// frontPool manages a pool of domain front candidates. Goroutine-safe.
// Take blocks until a front is available or ctx is cancelled.
// Return puts a front back after use. Replace atomically swaps the set.
type frontPool struct {
mu sync.Mutex
fronts []*front
ready chan *front
closed chan struct{}
}
func newFrontPool(readySize int) *frontPool {
if readySize <= 0 {
readySize = 500
}
return &frontPool{
ready: make(chan *front, readySize),
closed: make(chan struct{}),
}
}
// Replace atomically swaps the candidate set. Fronts that were previously
// known to be working (matched by providerID+domain+IP) retain their state.
func (p *frontPool) Replace(fronts []*front) {
p.mu.Lock()
defer p.mu.Unlock()
old := make(map[frontKey]*front, len(p.fronts))
for _, f := range p.fronts {
old[f.key()] = f
}
for _, f := range fronts {
if prev, ok := old[f.key()]; ok {
f.setLastSucceeded(prev.lastSucceededTime())
}
}
p.fronts = fronts
// Drain the ready channel; the crawler will repopulate it
for {
select {
case <-p.ready:
default:
return
}
}
}
// addReady marks a front as ready (working) and makes it available for Take.
func (p *frontPool) addReady(f *front) {
select {
case p.ready <- f:
default:
// ready channel full, drop
}
}
// readyCount returns the number of fronts in the ready queue.
func (p *frontPool) readyCount() int {
return len(p.ready)
}
// Take returns a working front, blocking until one is available or ctx is cancelled.
func (p *frontPool) Take(ctx context.Context) (*front, error) {
select {
case f := <-p.ready:
return f, nil
case <-ctx.Done():
return nil, ctx.Err()
case <-p.closed:
return nil, context.Canceled
}
}
// maxProviderScan bounds how many ready fronts TakePreferringUntried will pull
// while hunting for an untried provider, so a ready queue dominated by one
// provider can't turn a single take into a large drain-and-requeue.
const maxProviderScan = 8
// TakePreferringUntried returns a working front, preferring one whose ProviderID
// is not present in tried. It blocks (like Take) until at least one front is
// available, then scans a bounded number of additional ready fronts without
// blocking; every front it pulls but doesn't select is returned to the ready
// queue. If no untried provider is immediately available it falls back to the
// first front taken, so it never blocks longer than Take. A nil/empty tried map
// makes it behave exactly like Take.
func (p *frontPool) TakePreferringUntried(ctx context.Context, tried map[string]struct{}) (*front, error) {
first, err := p.Take(ctx)
if err != nil {
return nil, err
}
if _, seen := tried[first.ProviderID]; !seen {
return first, nil
}
skipped := []*front{first}
selected := first
scan := min(p.readyCount(), maxProviderScan)
for i := 0; i < scan; i++ {
f, ok := p.tryTakeReady()
if !ok {
break // nothing more ready right now
}
if _, seen := tried[f.ProviderID]; !seen {
selected = f
break
}
skipped = append(skipped, f)
}
for _, s := range skipped {
if s != selected {
p.addReady(s)
}
}
return selected, nil
}
// tryTakeReady returns a ready front without blocking, or ok=false when none is
// immediately available.
func (p *frontPool) tryTakeReady() (*front, bool) {
select {
case f := <-p.ready:
return f, true
default:
return nil, false
}
}
// Return puts a front back into the ready queue without updating its success
// timestamp. Use this when the front should be kept in rotation but no real
// round trip occurred (e.g. provider mapping miss). Pass requeue=false to
// mark the front as failed and remove it from rotation.
func (p *frontPool) Return(f *front, requeue bool) {
if requeue {
p.addReady(f)
} else {
f.markFailed()
}
}
// ReturnSuccess records a real successful round trip and puts the front back
// into the ready queue.
func (p *frontPool) ReturnSuccess(f *front) {
f.markSucceeded()
p.addReady(f)
}
// Close shuts down the pool, unblocking any pending Take calls.
func (p *frontPool) Close() {
select {
case <-p.closed:
// already closed
default:
close(p.closed)
}
}
// candidates returns a copy of all known fronts, with recently-succeeded
// fronts first and the rest shuffled.
func (p *frontPool) candidates() []*front {
p.mu.Lock()
n := len(p.fronts)
if n == 0 {
p.mu.Unlock()
return nil
}
// Pack pointer + int64 timestamp together (16 bytes vs 32 for time.Time).
// This halves the working-set size and makes sort comparisons cheaper
// (plain int64 compare vs time.Time.After with monotonic-clock logic).
type item struct {
f *front
ts int64 // UnixNano; 0 = never succeeded
}
items := make([]item, n)
for i, f := range p.fronts {
items[i].f = f
}
p.mu.Unlock()
// Snapshot timestamps outside the pool lock
for i := range items {
if t := items[i].f.lastSucceededTime(); !t.IsZero() {
items[i].ts = t.UnixNano()
}
}
// pdqsort via slices.SortFunc: faster than sort.Slice, no interface boxing
slices.SortFunc(items, func(a, b item) int {
if c := cmp.Compare(b.ts, a.ts); c != 0 {
return c // descending: most recent first
}
return cmp.Compare(a.f.IpAddress, b.f.IpAddress)
})
// Shuffle the non-succeeded tail
tail := 0
for tail < n && items[tail].ts != 0 {
tail++
}
rest := items[tail:]
rand.Shuffle(len(rest), func(i, j int) { rest[i], rest[j] = rest[j], rest[i] })
c := make([]*front, n)
for i := range items {
c[i] = items[i].f
}
return c
}