Skip to content

Commit 11496c4

Browse files
Moroka8WINGS-N
authored andcommitted
[client,server]: integrate WRAP per-packet obfuscation into relay datapath
Co-authored-by: Moroka8 <moroka8@mail.ru> (ported from cacggghp/vk-turn-proxy PR cacggghp#162, commit 0f24441; cipher selection extended to aes-ctr/chacha20-xor)
1 parent 57b411c commit 11496c4

9 files changed

Lines changed: 338 additions & 7 deletions

File tree

client/main.go

Lines changed: 105 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import (
2626
"time"
2727

2828
"github.com/cacggghp/vk-turn-proxy/internal/controlpath"
29+
"github.com/cacggghp/vk-turn-proxy/internal/wrap"
2930
"github.com/cacggghp/vk-turn-proxy/sessionproto"
3031
"github.com/cbeuw/connutil"
3132
"github.com/google/uuid"
@@ -1467,6 +1468,16 @@ type turnParams struct {
14671468
getCreds getCredsFunc
14681469
resolver *protectedResolver
14691470
credsManager *groupedCredsManager
1471+
// WRAP per-packet obfuscation. wrapCipher == WRAP_CIPHER_NONE /
1472+
// _UNSPECIFIED disables WRAP regardless of wrapKey.
1473+
wrapCipher sessionproto.WrapCipher
1474+
wrapKey []byte
1475+
// wrapMode controls fallback semantics when WRAP is configured:
1476+
// "off" — never wrap (wrapCipher should already be NONE)
1477+
// "preferred" — try WRAP, fall back to raw if no successfully unwrapped
1478+
// inbound packet arrives within wrapFallbackInboundTimeout
1479+
// "required" — fail hard if WRAP unwrap doesn't succeed; never fall back
1480+
wrapMode string
14701481
}
14711482

14721483
func oneTurnConnection(
@@ -1632,6 +1643,53 @@ func oneTurnConnection(
16321643
}
16331644
}
16341645
})
1646+
wrapCipher, err1 := wrap.New(turnParams.wrapCipher, turnParams.wrapKey)
1647+
if err1 != nil {
1648+
err = fmt.Errorf("WRAP cipher init failed: %s", err1)
1649+
return
1650+
}
1651+
wrapModeForAttempt := turnParams.wrapMode
1652+
// Fallback key is the peer (our server) address: WRAP support is a
1653+
// property of the peer, not of the VK TURN relay we route through.
1654+
fallbackKey := peer.String()
1655+
if wrapCipher != nil && wrapModeForAttempt == "preferred" && wrapDisabledForAddr(fallbackKey) {
1656+
log.Printf("[STREAM %d] WRAP marked unsupported for peer %s recently; using raw this attempt", streamID, fallbackKey)
1657+
wrapCipher = nil
1658+
}
1659+
if wrapCipher != nil {
1660+
log.Printf("[STREAM %d] WRAP active: cipher=%s mode=%s", streamID, turnParams.wrapCipher, wrapModeForAttempt)
1661+
}
1662+
var anyWrapInboundSuccess atomic.Bool
1663+
wrapActiveThisAttempt := wrapCipher != nil
1664+
// On return: if WRAP was active and no successful inbound unwrap
1665+
// happened (worker exited due to DTLS handshake timeout, peer
1666+
// closure, ctx cancellation, etc.), record peer as no-wrap so the
1667+
// maintain loop's next attempt goes raw within the TTL window.
1668+
defer func() {
1669+
if wrapActiveThisAttempt && wrapModeForAttempt == "preferred" && !anyWrapInboundSuccess.Load() {
1670+
markWrapDisabledForAddr(fallbackKey)
1671+
log.Printf("[STREAM %d] WRAP exit with no decoded inbound; disabling WRAP for peer %s", streamID, fallbackKey)
1672+
}
1673+
}()
1674+
if wrapActiveThisAttempt && wrapModeForAttempt == "preferred" {
1675+
go func() {
1676+
timer := time.NewTimer(wrapFallbackInboundTimeout)
1677+
defer timer.Stop()
1678+
select {
1679+
case <-turnctx.Done():
1680+
return
1681+
case <-timer.C:
1682+
if !anyWrapInboundSuccess.Load() {
1683+
markWrapDisabledForAddr(fallbackKey)
1684+
log.Printf(
1685+
"[STREAM %d] no WRAP-decoded inbound from peer %s in %s — disabling WRAP and reconnecting raw",
1686+
streamID, fallbackKey, wrapFallbackInboundTimeout,
1687+
)
1688+
turncancel()
1689+
}
1690+
}
1691+
}()
1692+
}
16351693
var addr atomic.Value
16361694
// Start read-loop on conn2 (output of DTLS)
16371695
go func() {
@@ -1654,7 +1712,16 @@ func oneTurnConnection(
16541712

16551713
addr.Store(addr1) // store peer
16561714

1657-
_, err1 = relayConn.WriteTo(buf[:n], peer)
1715+
payload := buf[:n]
1716+
if wrapCipher != nil {
1717+
sealed, sealErr := wrapCipher.Seal(payload)
1718+
if sealErr != nil {
1719+
log.Printf("[STREAM %d] WRAP seal failed: %s", streamID, sealErr)
1720+
return
1721+
}
1722+
payload = sealed
1723+
}
1724+
_, err1 = relayConn.WriteTo(payload, peer)
16581725
if err1 != nil {
16591726
if !shouldSuppressWorkerError(turnctx, err1) {
16601727
log.Printf("Failed: %s", err1)
@@ -1668,7 +1735,11 @@ func oneTurnConnection(
16681735
go func() {
16691736
defer wg.Done()
16701737
defer turncancel()
1671-
buf := make([]byte, 1600)
1738+
readBufLen := 1600
1739+
if wrapCipher != nil {
1740+
readBufLen += wrapCipher.Overhead()
1741+
}
1742+
buf := make([]byte, readBufLen)
16721743
for {
16731744
select {
16741745
case <-turnctx.Done():
@@ -1687,7 +1758,17 @@ func oneTurnConnection(
16871758
continue
16881759
}
16891760

1690-
_, err1 = conn2.WriteTo(buf[:n], addr1)
1761+
payload := buf[:n]
1762+
if wrapCipher != nil {
1763+
plain, openErr := wrapCipher.Open(payload)
1764+
if openErr != nil {
1765+
log.Printf("[STREAM %d] WRAP unwrap failed (%d bytes): %s", streamID, n, openErr)
1766+
continue
1767+
}
1768+
anyWrapInboundSuccess.Store(true)
1769+
payload = plain
1770+
}
1771+
_, err1 = conn2.WriteTo(payload, addr1)
16911772
if err1 != nil {
16921773
if !shouldSuppressWorkerError(turnctx, err1) {
16931774
log.Printf("Failed: %s", err1)
@@ -2016,6 +2097,10 @@ func main() { //nolint:cyclop
20162097
setStrategy := func(_ sessionproto.Mode, _ uint32, _ int) {}
20172098
return unifiedGetCreds, setStrategy
20182099
}
2100+
wrapCipherSel, wrapKey, wrapMode, err := resolveWrapConfig(opts.wrapMode, opts.wrapCipher, opts.wrapKeyHex)
2101+
if err != nil {
2102+
log.Panicf("WRAP config: %v", err)
2103+
}
20192104
params := &turnParams{
20202105
host: opts.host,
20212106
port: opts.port,
@@ -2024,6 +2109,9 @@ func main() { //nolint:cyclop
20242109
getCreds: nil,
20252110
resolver: peerResolver,
20262111
credsManager: vkLinkManager,
2112+
wrapCipher: wrapCipherSel,
2113+
wrapKey: wrapKey,
2114+
wrapMode: wrapMode,
20272115
}
20282116
sessionID := []byte(nil)
20292117

@@ -2271,10 +2359,24 @@ func main() { //nolint:cyclop
22712359
false,
22722360
)
22732361
if !waitForReady(ctx, okchan, mainlineBootstrapTimeout) {
2362+
// If WRAP was attempted, the watchdog in each worker has
2363+
// by now marked the peer as no-wrap; pre-emptively mark
2364+
// it here too so any racing worker also goes raw. Give
2365+
// bootstrap one more shot before tearing the process
2366+
// down so Android does not enter a respawn loop that
2367+
// wipes the wrap-disabled cache each time.
2368+
if params.wrapMode == "preferred" {
2369+
markWrapDisabledForAddr(peer.String())
2370+
log.Printf("bootstrap timed out; forcing WRAP fallback for peer %s and retrying", peer.String())
2371+
if waitForReady(ctx, okchan, mainlineBootstrapTimeout) {
2372+
goto mainlineBootstrapDone
2373+
}
2374+
}
22742375
runtimeCancel()
22752376
runtimeWG.Wait()
22762377
log.Fatalf("failed to bootstrap mainline session")
22772378
}
2379+
mainlineBootstrapDone:
22782380

22792381
supportedVersion := waitForProbeVersion(ctx, probeResult, muProbeTimeout)
22802382
activeMainlineControl := waitForMainlineControlHandle(ctx, mainlineControl, muProbeTimeout)

client/options.go

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,10 @@ type clientOptions struct {
5050
roomExchangeDisplayName string
5151
roomExchangeE2EEnabled bool
5252
roomExchangeE2ESecret string
53+
54+
wrapMode string // off|preferred|required (preferred ≡ required without negotiation in v1)
55+
wrapCipher string // aes-ctr|chacha20-xor
56+
wrapKeyHex string // 64-char hex (32 bytes); empty + mode!=off generates fresh key
5357
}
5458

5559
func newClientFlagSet(program string, output io.Writer) (*flag.FlagSet, *clientOptions) {
@@ -93,6 +97,9 @@ func newClientFlagSet(program string, output io.Writer) (*flag.FlagSet, *clientO
9397
fs.StringVar(&opts.roomExchangeDisplayName, "room-exchange-display-name", "", "display name delivered alongside the room id in the room-exchange handshake")
9498
fs.BoolVar(&opts.roomExchangeE2EEnabled, "room-exchange-e2e-enabled", false, "advertise that wb-stream traffic will be E2E-encrypted")
9599
fs.StringVar(&opts.roomExchangeE2ESecret, "room-exchange-e2e-secret", "", "optional base64-encoded E2E secret to share with the server")
100+
fs.StringVar(&opts.wrapMode, "wrap-mode", "off", "WRAP per-packet obfuscation mode: off|preferred|required")
101+
fs.StringVar(&opts.wrapCipher, "wrap-cipher", "aes-ctr", "WRAP cipher: aes-ctr|chacha20-xor")
102+
fs.StringVar(&opts.wrapKeyHex, "wrap-key", "", "WRAP key (64-char hex, 32 bytes); empty + mode!=off generates a fresh key")
96103
fs.Usage = func() {
97104
cliutil.Fprintf(fs.Output(), "Usage:\n %s -peer <host:port> -vk-link <link> [flags]\n %s -peer <host:port> -yandex-link <link> [flags]\n\n", program, program)
98105
cliutil.Fprintln(fs.Output(), "Examples:")
@@ -122,6 +129,19 @@ func parseClientOptions(args []string, program string, stdout, stderr io.Writer)
122129
if opts.credsGroupSize < 1 {
123130
opts.credsGroupSize = 1
124131
}
132+
opts.wrapMode = strings.ToLower(strings.TrimSpace(opts.wrapMode))
133+
switch opts.wrapMode {
134+
case "off", "preferred", "required":
135+
default:
136+
opts.wrapMode = "off"
137+
}
138+
opts.wrapCipher = strings.ToLower(strings.TrimSpace(opts.wrapCipher))
139+
switch opts.wrapCipher {
140+
case "aes-ctr", "chacha20-xor":
141+
default:
142+
opts.wrapCipher = "aes-ctr"
143+
}
144+
opts.wrapKeyHex = strings.TrimSpace(opts.wrapKeyHex)
125145

126146
opts.wbStreamRoomID = strings.TrimSpace(opts.wbStreamRoomID)
127147
opts.wbStreamRoomIDs = strings.TrimSpace(opts.wbStreamRoomIDs)

client/wrap_config.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
package main
2+
3+
import (
4+
"crypto/rand"
5+
"encoding/hex"
6+
"fmt"
7+
"log"
8+
9+
"github.com/cacggghp/vk-turn-proxy/internal/wrap"
10+
"github.com/cacggghp/vk-turn-proxy/sessionproto"
11+
)
12+
13+
// resolveWrapConfig validates and normalizes the WRAP CLI options into a
14+
// (cipher, key, mode) triple consumable by turnParams. Returns NONE/nil
15+
// when mode is "off". Generates a random key if mode is enabled and key
16+
// is empty.
17+
func resolveWrapConfig(mode, cipherStr, keyHex string) (sessionproto.WrapCipher, []byte, string, error) {
18+
if mode == "off" || mode == "" {
19+
return sessionproto.WrapCipher_WRAP_CIPHER_NONE, nil, "off", nil
20+
}
21+
22+
var selected sessionproto.WrapCipher
23+
switch cipherStr {
24+
case "aes-ctr":
25+
selected = sessionproto.WrapCipher_WRAP_CIPHER_AES_256_CTR
26+
case "chacha20-xor":
27+
selected = sessionproto.WrapCipher_WRAP_CIPHER_CHACHA20_XOR
28+
default:
29+
return 0, nil, "", fmt.Errorf("unsupported -wrap-cipher %q", cipherStr)
30+
}
31+
32+
var key []byte
33+
if keyHex == "" {
34+
gen := make([]byte, wrap.KeyLen)
35+
if _, err := rand.Read(gen); err != nil {
36+
return 0, nil, "", fmt.Errorf("WRAP key generation: %w", err)
37+
}
38+
key = gen
39+
log.Printf("WRAP: auto-generated key (set -wrap-key=%s on both ends to fix it)", hex.EncodeToString(key))
40+
} else {
41+
decoded, err := hex.DecodeString(keyHex)
42+
if err != nil {
43+
return 0, nil, "", fmt.Errorf("-wrap-key invalid hex: %w", err)
44+
}
45+
if len(decoded) != wrap.KeyLen {
46+
return 0, nil, "", fmt.Errorf("-wrap-key must decode to %d bytes (got %d)", wrap.KeyLen, len(decoded))
47+
}
48+
key = decoded
49+
}
50+
return selected, key, mode, nil
51+
}

client/wrap_fallback.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
package main
2+
3+
import (
4+
"sync"
5+
"time"
6+
)
7+
8+
// Client-side WRAP fallback heuristic.
9+
//
10+
// WRAP wraps every datagram travelling between client and TURN relay
11+
// endpoint. DTLS handshake runs on top of (un)wrapped bytes, so a server
12+
// that does not speak WRAP simply cannot parse our wrapped ClientHello
13+
// and silently drops it — there is no in-band way to "downgrade" mid-DTLS.
14+
//
15+
// For mode "preferred" we therefore implement a heuristic: spawn a
16+
// watchdog when WRAP is active. If no datagram from the TURN relay is
17+
// successfully unwrapped within wrapFallbackInboundTimeout, mark the
18+
// server address as "WRAP unsupported" for wrapFallbackAddrTTL and tear
19+
// the worker down. The maintain loop reconnects and the next attempt
20+
// (within the TTL window) skips WRAP for that address.
21+
22+
const (
23+
// 5s gives DTLS time for ~3 retransmits (1s, 2s, 4s) — long enough to
24+
// be sure the peer is not replying, short enough that two attempts
25+
// (WRAP then raw) comfortably fit inside the 30s mainline bootstrap
26+
// budget even when TURN allocate is slow.
27+
wrapFallbackInboundTimeout = 5 * time.Second
28+
wrapFallbackAddrTTL = 5 * time.Minute
29+
)
30+
31+
var wrapDisabledAddrs sync.Map // map[string]time.Time
32+
33+
func wrapDisabledForAddr(addr string) bool {
34+
v, ok := wrapDisabledAddrs.Load(addr)
35+
if !ok {
36+
return false
37+
}
38+
at, ok := v.(time.Time)
39+
if !ok || time.Since(at) > wrapFallbackAddrTTL {
40+
wrapDisabledAddrs.Delete(addr)
41+
return false
42+
}
43+
return true
44+
}
45+
46+
func markWrapDisabledForAddr(addr string) {
47+
if addr == "" {
48+
return
49+
}
50+
wrapDisabledAddrs.Store(addr, time.Now())
51+
}

internal/wrap/conn.go

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@ package wrap
33
import (
44
"log"
55
"net"
6+
7+
dtlsnet "github.com/pion/dtls/v3/pkg/net"
68
)
79

810
// readBufSize bounds how much we read from the underlying conn per packet.
@@ -61,3 +63,31 @@ func (w *wrappedConn) WriteTo(p []byte, addr net.Addr) (int, error) {
6163
}
6264
return len(p), nil
6365
}
66+
67+
// PacketListener wraps a dtls.PacketListener so every accepted connection has
68+
// the supplied Cipher applied to its reads/writes.
69+
//
70+
// Returning the inner listener verbatim when cipher is nil keeps call-sites
71+
// symmetric for the "no obfuscation negotiated" path.
72+
func PacketListener(inner dtlsnet.PacketListener, c Cipher) dtlsnet.PacketListener {
73+
if c == nil {
74+
return inner
75+
}
76+
return &wrappedListener{inner: inner, cipher: c}
77+
}
78+
79+
type wrappedListener struct {
80+
inner dtlsnet.PacketListener
81+
cipher Cipher
82+
}
83+
84+
func (l *wrappedListener) Accept() (net.PacketConn, net.Addr, error) {
85+
pc, addr, err := l.inner.Accept()
86+
if err != nil {
87+
return pc, addr, err
88+
}
89+
return PacketConn(pc, l.cipher), addr, nil
90+
}
91+
92+
func (l *wrappedListener) Close() error { return l.inner.Close() }
93+
func (l *wrappedListener) Addr() net.Addr { return l.inner.Addr() }

internal/wrap/conn_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,14 +37,14 @@ func TestPacketConnRoundTrip(t *testing.T) {
3737
ws := PacketConn(server, serverCipher)
3838

3939
payload := []byte("vk-turn obfuscation packet conn round-trip — should arrive intact")
40-
if err := wc.SetWriteDeadline(time.Now().Add(2 * time.Second)); err != nil {
40+
if err = wc.SetWriteDeadline(time.Now().Add(2 * time.Second)); err != nil {
4141
t.Fatal(err)
4242
}
43-
if _, err := wc.WriteTo(payload, server.LocalAddr()); err != nil {
43+
if _, err = wc.WriteTo(payload, server.LocalAddr()); err != nil {
4444
t.Fatalf("WriteTo: %v", err)
4545
}
4646

47-
if err := ws.SetReadDeadline(time.Now().Add(2 * time.Second)); err != nil {
47+
if err = ws.SetReadDeadline(time.Now().Add(2 * time.Second)); err != nil {
4848
t.Fatal(err)
4949
}
5050
buf := make([]byte, 4096)

0 commit comments

Comments
 (0)