Skip to content

Commit 6ea4ae8

Browse files
Moroka8WINGS-N
authored andcommitted
[client]: cache all VK TURN addresses and rotate per worker
Co-authored-by: Moroka8 <moroka8@mail.ru> (ported from cacggghp/vk-turn-proxy PR cacggghp#162, commit 272aa36)
1 parent d42d42d commit 6ea4ae8

4 files changed

Lines changed: 166 additions & 40 deletions

File tree

client/group_creds.go

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -91,7 +91,15 @@ func (m *groupedCredsManager) GetCredsForWorker(workerID int) (string, string, s
9191
if err != nil {
9292
return "", "", "", err
9393
}
94-
return cred.user, cred.pass, cred.addr, nil
94+
return cred.user, cred.pass, pickStreamServerAddr(workerID, cred.addrs), nil
95+
}
96+
97+
// ReportSetupFailure advances the rotation offset for the worker and marks the
98+
// failing TURN address as cooling-down, so subsequent GetCredsForWorker calls
99+
// hand out a different endpoint within the same credential.
100+
func (m *groupedCredsManager) ReportSetupFailure(workerID int, addr string) {
101+
rotateStreamServer(workerID)
102+
markTURNServerCooldown(addr)
95103
}
96104

97105
func (m *groupedCredsManager) ReportWorkerError(workerID int, err error) {

client/main.go

Lines changed: 44 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -340,18 +340,18 @@ func vkDelayRandom(minMs, maxMs int) {
340340
time.Sleep(time.Duration(ms) * time.Millisecond)
341341
}
342342

343-
func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInteractiveFallback bool) (string, string, string, time.Duration, error) {
343+
func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInteractiveFallback bool) (string, string, []string, time.Duration, error) {
344344
if remaining := captchaLockoutRemaining(); remaining > 0 {
345345
emitCaptchaLockoutStatus(remaining)
346-
return "", "", "", 0, fmt.Errorf("CAPTCHA_WAIT_REQUIRED: global lockout active for %s", remaining.Round(time.Second))
346+
return "", "", nil, 0, fmt.Errorf("CAPTCHA_WAIT_REQUIRED: global lockout active for %s", remaining.Round(time.Second))
347347
}
348348

349349
profile := getRandomProfile()
350350
name := generateName()
351351
escapedName := neturl.QueryEscape(name)
352352
client, err := resolver.newTLSHTTPClient(profile, 20*time.Second)
353353
if err != nil {
354-
return "", "", "", 0, fmt.Errorf("failed to initialize tls client: %w", err)
354+
return "", "", nil, 0, fmt.Errorf("failed to initialize tls client: %w", err)
355355
}
356356
defer client.CloseIdleConnections()
357357

@@ -410,16 +410,16 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
410410

411411
resp, err = doRequest(data, url)
412412
if err != nil {
413-
return "", "", "", 0, fmt.Errorf("request error:%s", err)
413+
return "", "", nil, 0, fmt.Errorf("request error:%s", err)
414414
}
415415

416416
dataMap, ok := resp["data"].(map[string]interface{})
417417
if !ok {
418-
return "", "", "", 0, fmt.Errorf("unexpected anon token response: %v", resp)
418+
return "", "", nil, 0, fmt.Errorf("unexpected anon token response: %v", resp)
419419
}
420420
token1, ok := dataMap["access_token"].(string)
421421
if !ok {
422-
return "", "", "", 0, fmt.Errorf("missing access_token in response: %v", resp)
422+
return "", "", nil, 0, fmt.Errorf("missing access_token in response: %v", resp)
423423
}
424424

425425
vkDelayRandom(100, 150)
@@ -440,14 +440,14 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
440440
for attempt := 0; attempt <= maxCaptchaAttempts; attempt++ {
441441
resp, err = doRequest(data, url)
442442
if err != nil {
443-
return "", "", "", 0, fmt.Errorf("request error:%s", err)
443+
return "", "", nil, 0, fmt.Errorf("request error:%s", err)
444444
}
445445

446446
if errObj, hasErr := resp["error"].(map[string]interface{}); hasErr {
447447
errCode, _ := errObj["error_code"].(float64)
448448
if errCode == 14 {
449449
if attempt == maxCaptchaAttempts {
450-
return "", "", "", 0, wrapCaptchaFailure(fmt.Errorf("captcha failed after %d attempts", maxCaptchaAttempts), allowInteractiveFallback)
450+
return "", "", nil, 0, wrapCaptchaFailure(fmt.Errorf("captcha failed after %d attempts", maxCaptchaAttempts), allowInteractiveFallback)
451451
}
452452

453453
captchaErr := parseVkCaptchaError(errObj)
@@ -507,7 +507,7 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
507507
}
508508
}
509509
if solveErr != nil {
510-
return "", "", "", 0, wrapCaptchaFailure(fmt.Errorf("smart captcha solve error: %w", solveErr), allowInteractiveFallback)
510+
return "", "", nil, 0, wrapCaptchaFailure(fmt.Errorf("smart captcha solve error: %w", solveErr), allowInteractiveFallback)
511511
}
512512
storeCachedCaptchaToken(successToken)
513513
captchaAttempt := captchaErr.CaptchaAttempt
@@ -538,7 +538,7 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
538538
profile.UserAgent,
539539
)
540540
if solveErr != nil {
541-
return "", "", "", 0, wrapCaptchaFailure(fmt.Errorf("captcha solve error: %w", solveErr), false)
541+
return "", "", nil, 0, wrapCaptchaFailure(fmt.Errorf("captcha solve error: %w", solveErr), false)
542542
}
543543
data = fmt.Sprintf(
544544
"vk_join_link=https://vk.com/call/join/%s&name=%s&access_token=%s&captcha_sid=%s&captcha_key=%s",
@@ -562,7 +562,7 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
562562
profile.UserAgent,
563563
)
564564
if solveErr != nil {
565-
return "", "", "", 0, wrapCaptchaFailure(fmt.Errorf("captcha solve error: %w", solveErr), true)
565+
return "", "", nil, 0, wrapCaptchaFailure(fmt.Errorf("captcha solve error: %w", solveErr), true)
566566
}
567567
data = fmt.Sprintf(
568568
"vk_join_link=https://vk.com/call/join/%s&name=%s&access_token=%s&captcha_sid=%s&captcha_key=%s",
@@ -575,16 +575,16 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
575575
}
576576
continue
577577
}
578-
return "", "", "", 0, fmt.Errorf("VK API error: %v", errObj)
578+
return "", "", nil, 0, fmt.Errorf("VK API error: %v", errObj)
579579
}
580580

581581
responseMap, ok := resp["response"].(map[string]interface{})
582582
if !ok {
583-
return "", "", "", 0, fmt.Errorf("unexpected getAnonymousToken response: %v", resp)
583+
return "", "", nil, 0, fmt.Errorf("unexpected getAnonymousToken response: %v", resp)
584584
}
585585
token2, ok = responseMap["token"].(string)
586586
if !ok {
587-
return "", "", "", 0, fmt.Errorf("missing token in response: %v", resp)
587+
return "", "", nil, 0, fmt.Errorf("missing token in response: %v", resp)
588588
}
589589
if usedAutoCaptcha {
590590
log.Printf("VK smart captcha accepted by auth endpoint")
@@ -598,7 +598,7 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
598598

599599
resp, err = doRequest(data, url)
600600
if err != nil {
601-
return "", "", "", 0, fmt.Errorf("request error:%s", err)
601+
return "", "", nil, 0, fmt.Errorf("request error:%s", err)
602602
}
603603

604604
token3 := resp["session_key"].(string)
@@ -609,13 +609,16 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
609609

610610
resp, err = doRequest(data, url)
611611
if err != nil {
612-
return "", "", "", 0, fmt.Errorf("request error:%s", err)
612+
return "", "", nil, 0, fmt.Errorf("request error:%s", err)
613613
}
614614

615615
turnServer := resp["turn_server"].(map[string]interface{})
616616
user := turnServer["username"].(string)
617617
pass := turnServer["credential"].(string)
618-
turn := turnServer["urls"].([]interface{})[0].(string)
618+
urlsRaw, ok := turnServer["urls"].([]interface{})
619+
if !ok || len(urlsRaw) == 0 {
620+
return "", "", nil, 0, fmt.Errorf("missing or empty urls in turn_server")
621+
}
619622

620623
var lifetime time.Duration
621624
if rawLifetime, ok := turnServer["lifetime"].(float64); ok && rawLifetime > 0 {
@@ -624,10 +627,22 @@ func getVkCredsWithFallback(link string, resolver *protectedResolver, allowInter
624627
lifetime = time.Duration(rawTTL) * time.Second
625628
}
626629

627-
clean := strings.Split(turn, "?")[0]
628-
address := strings.TrimPrefix(strings.TrimPrefix(clean, "turn:"), "turns:")
630+
var addresses []string
631+
for _, u := range urlsRaw {
632+
urlStr, ok := u.(string)
633+
if !ok {
634+
continue
635+
}
636+
clean := strings.Split(urlStr, "?")[0]
637+
address := strings.TrimPrefix(strings.TrimPrefix(clean, "turn:"), "turns:")
638+
addresses = append(addresses, address)
639+
}
640+
if len(addresses) == 0 {
641+
return "", "", nil, 0, fmt.Errorf("no valid TURN addresses parsed from urls")
642+
}
643+
log.Printf("VK Auth: TURN urls (%d) %v", len(addresses), addresses)
629644

630-
return user, pass, address, lifetime, nil
645+
return user, pass, addresses, lifetime, nil
631646
}
632647

633648
func getYandexCreds(link string, resolver *protectedResolver) (string, string, string, error) {
@@ -1941,11 +1956,11 @@ func main() { //nolint:cyclop
19411956
credsGroupSize := max(1, opts.credsGroupSize)
19421957
numGroups := max(1, ceilDiv(opts.n, credsGroupSize))
19431958
vkFetch := func(fctx context.Context, hash string, allowInteractive bool) (turnCred, error) {
1944-
user, pass, addr, lifetime, err := getVkCredsWithFallback(hash, peerResolver, allowInteractive)
1959+
user, pass, addrs, lifetime, err := getVkCredsWithFallback(hash, peerResolver, allowInteractive)
19451960
if err != nil {
19461961
return turnCred{}, err
19471962
}
1948-
return turnCred{user: user, pass: pass, addr: addr, lifetime: lifetime}, nil
1963+
return turnCred{user: user, pass: pass, addrs: addrs, lifetime: lifetime}, nil
19491964
}
19501965
vkLinkManager = newGroupedCredsManager(ctx, numGroups, credsGroupSize, tracker, vkFetch)
19511966
log.Printf(
@@ -1967,14 +1982,17 @@ func main() { //nolint:cyclop
19671982
if opts.n <= 0 {
19681983
opts.n = 1
19691984
}
1970-
yandexBase := func(s string, _ bool) (string, string, string, time.Duration, error) {
1985+
yandexBase := func(s string, _ bool) (string, string, []string, time.Duration, error) {
19711986
user, pass, addr, err := getYandexCreds(s, peerResolver)
1972-
return user, pass, addr, 0, err
1987+
if err != nil {
1988+
return "", "", nil, 0, err
1989+
}
1990+
return user, pass, []string{addr}, 0, nil
19731991
}
19741992
yandexPool := poolCreds(yandexBase, 1)
19751993
yandexLink := link
1976-
unifiedGetCreds = func(_ int) (string, string, string, error) {
1977-
return yandexPool(yandexLink)
1994+
unifiedGetCreds = func(workerID int) (string, string, string, error) {
1995+
return yandexPool(yandexLink, workerID)
19781996
}
19791997
}
19801998
configuredPoolSize := max(1, opts.n)

client/pool_creds.go

Lines changed: 22 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -11,15 +11,24 @@ import (
1111
type turnCred struct {
1212
user string
1313
pass string
14-
addr string
14+
addrs []string
1515
bornAt time.Time
1616
lifetime time.Duration
1717
isSecondaryLink bool
1818
}
1919

20-
type pooledGetCredsFunc func(string, bool) (string, string, string, time.Duration, error)
20+
// primaryAddr returns the canonical address for identity comparison; an empty
21+
// addrs slice yields "" so dedup/log paths stay safe.
22+
func (cred turnCred) primaryAddr() string {
23+
if len(cred.addrs) == 0 {
24+
return ""
25+
}
26+
return cred.addrs[0]
27+
}
28+
29+
type pooledGetCredsFunc func(string, bool) (string, string, []string, time.Duration, error)
2130

22-
type pooledGetCredsResult func(link string) (string, string, string, error)
31+
type pooledGetCredsResult func(link string, workerID int) (string, string, string, error)
2332

2433
type adaptivePoolConfig struct {
2534
minSize int
@@ -165,7 +174,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
165174

166175
appendIfNewLocked := func(cred turnCred) bool {
167176
for _, existing := range state.pool {
168-
if existing.user == cred.user && existing.pass == cred.pass && existing.addr == cred.addr {
177+
if existing.user == cred.user && existing.pass == cred.pass && existing.primaryAddr() == cred.primaryAddr() {
169178
return false
170179
}
171180
}
@@ -196,7 +205,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
196205
state.mu.Unlock()
197206
}()
198207

199-
user, pass, addr, lifetime, err := f(link, false)
208+
user, pass, addrs, lifetime, err := f(link, false)
200209
if err != nil {
201210
state.mu.Lock()
202211
state.backgroundRetryAfter = time.Now().Add(backgroundPoolRetryCooldown)
@@ -213,7 +222,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
213222
added := appendIfNewLocked(turnCred{
214223
user: user,
215224
pass: pass,
216-
addr: addr,
225+
addrs: addrs,
217226
bornAt: time.Now(),
218227
lifetime: lifetime,
219228
})
@@ -229,7 +238,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
229238
}()
230239
}
231240

232-
return func(link string) (string, string, string, error) {
241+
return func(link string, workerID int) (string, string, string, error) {
233242
for {
234243
state.mu.Lock()
235244
expireIfNeededLocked()
@@ -244,7 +253,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
244253
if currentSize < desiredSize {
245254
startBackgroundFill(link)
246255
}
247-
return cred.user, cred.pass, cred.addr, nil
256+
return cred.user, cred.pass, pickStreamServerAddr(workerID, cred.addrs), nil
248257
}
249258

250259
if state.foregroundFillRunning {
@@ -259,7 +268,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
259268
if currentSize < desiredSize {
260269
startBackgroundFill(link)
261270
}
262-
return cred.user, cred.pass, cred.addr, nil
271+
return cred.user, cred.pass, pickStreamServerAddr(workerID, cred.addrs), nil
263272
}
264273
state.mu.Unlock()
265274
continue
@@ -268,15 +277,15 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
268277
state.foregroundFillRunning = true
269278
state.mu.Unlock()
270279

271-
user, pass, addr, lifetime, err := f(link, true)
280+
user, pass, addrs, lifetime, err := f(link, true)
272281

273282
state.mu.Lock()
274283
state.foregroundFillRunning = false
275284
if err == nil {
276285
_ = appendIfNewLocked(turnCred{
277286
user: user,
278287
pass: pass,
279-
addr: addr,
288+
addrs: addrs,
280289
bornAt: time.Now(),
281290
lifetime: lifetime,
282291
})
@@ -291,7 +300,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
291300
if currentSize < desiredSize {
292301
startBackgroundFill(link)
293302
}
294-
return user, pass, addr, nil
303+
return user, pass, pickStreamServerAddr(workerID, addrs), nil
295304
}
296305

297306
state.cond.Broadcast()
@@ -304,7 +313,7 @@ func poolCredsDynamic(f pooledGetCredsFunc, targetPoolSize func() int) pooledGet
304313
if currentSize < desiredSize {
305314
startBackgroundFill(link)
306315
}
307-
return cred.user, cred.pass, cred.addr, nil
316+
return cred.user, cred.pass, pickStreamServerAddr(workerID, cred.addrs), nil
308317
}
309318
state.mu.Unlock()
310319

0 commit comments

Comments
 (0)