Files
EasyTier/easytier-go/internal/hostabi/socket_wire.go
T
KKRainbow e313ba8efb fix(wasi): reduce startup RSS, align ABI v4 and wire v3 (#2588)
* feat(easytier-go): accept data plane ABI v4

The core raised DATA_PLANE_ABI_VERSION to 4 in f26c2aa1 ("feat(wasi):
run EasyTier core on Cloudflare Workers and browsers") to advertise the
new guest exports easytier_data_plane_tcp_shutdown_write_submit and
easytier_data_plane_tcp_shutdown_write_result_take, which let a host
half-close guest TCP streams.

The Go host has no caller for half-close: net.Conn exposes only Close,
which already tears down both directions, so the v3 behavior is
preserved and no new plumbing is added. Without this bump, any artifact
rebuilt from current core source is rejected at instance creation with
"unsupported EasyTier data plane ABI version 4, want 3".

* fix(easytier-go): decode socket options wire v3, rebuild core artifact

38e2a621 ("refactor(ohos): 拆分 OHRS 包并按 socket 精细保护 VPN 流量")
raised the host socket options wire format from version 2 to 3: a
need_protect byte is appended after the purpose byte in TCP connect,
UDP bind, and TCP listen options, shifting bind_device one byte later.
The Go hostABI decoders were never updated, so every socket operation
from a HEAD-built core was rejected as "invalid options" and required
listeners failed to start.

Bump the accepted wire version to 3 and skip the need_protect byte in
all three decoders. The byte requests VPN socket protection, which only
Android-style VPN hosts can honor; on every other platform sockets are
already protected, so reading and discarding it is correct.

Regenerate the embedded core artifact and protobuf bindings from HEAD
(599e4eac) so the shipped artifact matches the Go host again. The proto
regeneration also picks up schema fields added since the last embed
(e.g. prefer_peer_relay).

* perf(easytier-go): release compiler garbage after host init

wazero's optimizing compiler allocates ~100MB of throwaway state on
the Go heap while compiling the embedded core. Go's runtime does not
return that memory to the OS after the initiating GC, so the process
retained the compilation peak for its entire lifetime: RSS sat at
~164MB before any instance or network activity.

Call debug.FreeOSMemory() once after the module is instantiated.
NewHost is a one-time initialization path outside the dataplane, so
the stop-the-world pass is safe here. Measured RSS after host.New
drops from ~164MB to ~64MB; no behavioral change.

* chore(easytier-js): bump toolchain dependencies

- vitest 2.1.9 -> 3.2.7 (all three packages)
- esbuild 0.25.9 -> 0.28.2 (browser bundler)
- wrangler 4.114.0 -> 4.134.0 (cloudflare + web example)
- @cloudflare/workers-types 5.20260724.1 -> 5.20260917.1
- binaryen 131.0.0 -> 132.0.0 (JS bindings only; the wasm build
  uses the standalone wasm-opt binary fetched by the build script)

Not bumped: typescript stays at 5.9.3 (latest 5.x; 7.0 is a major
jump not worth taking for this workspace), vite stays at 5.4.21
(web-example only; 5->8 spans three majors).

Verified with pnpm test (28 + 2 + 2 tests across runtime, browser,
cloudflare) and pnpm check (tsc + wrangler deploy --dry-run).
pnpm-workspace.yaml gained three minimumReleaseAgeExclude entries
recorded automatically by pnpm for freshly published versions.

* fix(easytier-js): align host wire format with core wire v3

The core socket options wire format moved to version 3 in 38e2a621
(need_protect byte after purpose). The JS websocket-host still
required version 2 and rejected every TCP bind from a HEAD-built core,
breaking browser port leases.

- websocket-host.ts: accept version 3, minimum length 49. The
  need_protect byte sits after purpose (offset 43) and is a no-op
  outside Android VPN hosts, so the decoder just skips it.
- websocket-host.test.ts: update the test encoder to emit v3.
- binaryen stays at 131.0.0 to match script/build-wasi-core.sh, which
  intentionally pins binaryen 131 for the Go-side embedded core.
  The npm binaryen package provides the wasm-opt binary used by
  build-wasm.mjs, so keeping both sides on the same version avoids
  divergent optimization output.
- pnpm-workspace.yaml: drop two stale minimumReleaseAgeExclude entries
  for @cloudflare/workers-types versions no longer in the lockfile.

* test(easytier-go): cover socket options wire v3 decoders

Direct unit tests for decodeTCPConnectOptions, decodeUDPBindOptions, and
decodeTCPListenOptions with wire v3 documents. Covers combinations of
socket mark, netns, bind device, and local address presence, plus
rejection of wire v2.

Resolves review feedback on PR #2588.
2026-09-19 00:13:47 +08:00

436 lines
12 KiB
Go

package hostabi
import (
"encoding/binary"
"fmt"
"net"
"unicode/utf8"
"github.com/easytier/easytier/easytier-go/platform"
"github.com/metacubex/wazero/api"
)
const (
socketAddressLen = 27
tcpSocketResultLen = 62
boundSocketResultLen = 35
maxFactoryOptionsSize = 4096
)
func readOwnedOptions(module api.Module, pointer, length uint32) ([]byte, bool) {
if length > maxFactoryOptionsSize {
return nil, false
}
options, ok := module.Memory().Read(pointer, length)
if !ok {
return nil, false
}
return append([]byte(nil), options...), true
}
// optionsWireVersion matches OPTIONS_VERSION in
// easytier-core/src/wasi/wire/options.rs. Version 3 appends a need_protect
// byte after the purpose byte in every socket options document. The byte
// requests VPN socket protection, which only Android-style VPN hosts can
// honor; every other platform treats sockets as already protected, so the
// decoders read and discard it.
const optionsWireVersion = 3
func decodeTCPConnectOptions(encoded []byte) (platform.TCPConnectOptions, error) {
if len(encoded) < 76 || encoded[0] != optionsWireVersion {
return platform.TCPConnectOptions{}, fmt.Errorf("invalid TCP connect options")
}
remote, err := decodeSocketAddress(encoded[1:28], false)
if err != nil || remote == nil {
return platform.TCPConnectOptions{}, fmt.Errorf("invalid TCP remote address")
}
local, err := decodeSocketAddress(encoded[28:55], true)
if err != nil {
return platform.TCPConnectOptions{}, err
}
socketContext, remainder, err := decodeSocketContext(encoded[55:])
if err != nil {
return platform.TCPConnectOptions{}, fmt.Errorf("invalid TCP socket context: %w", err)
}
if len(remainder) < 10 {
return platform.TCPConnectOptions{}, fmt.Errorf("truncated TCP bind policy")
}
bind, err := decodeTCPBindPolicy(
udpToTCPAddr(local),
socketContext,
remainder[0],
remainder[1],
remainder[2],
remainder[5:],
)
if err != nil {
return platform.TCPConnectOptions{}, err
}
purpose, err := decodeTCPConnectPurpose(remainder[3])
if err != nil {
return platform.TCPConnectOptions{}, err
}
return platform.TCPConnectOptions{
RemoteAddr: &net.TCPAddr{IP: remote.IP, Port: remote.Port, Zone: remote.Zone},
Bind: bind,
Purpose: purpose,
}, nil
}
func decodeUDPBindOptions(encoded []byte) (platform.UDPBindOptions, error) {
if len(encoded) < 49 || encoded[0] != optionsWireVersion {
return platform.UDPBindOptions{}, fmt.Errorf("invalid UDP bind options")
}
local, err := decodeSocketAddress(encoded[1:28], true)
if err != nil {
return platform.UDPBindOptions{}, err
}
socketContext, remainder, err := decodeSocketContext(encoded[28:])
if err != nil {
return platform.UDPBindOptions{}, fmt.Errorf("invalid UDP socket context: %w", err)
}
if len(remainder) < 10 {
return platform.UDPBindOptions{}, fmt.Errorf("truncated UDP bind policy")
}
reuseAddr, err := decodeWireBool("UDP reuse_addr", remainder[0])
if err != nil {
return platform.UDPBindOptions{}, err
}
reusePort, err := decodeWireBool("UDP reuse_port", remainder[1])
if err != nil {
return platform.UDPBindOptions{}, err
}
onlyV6, err := decodeWireBool("UDP only_v6", remainder[2])
if err != nil {
return platform.UDPBindOptions{}, err
}
purpose, err := decodeUDPBindPurpose(remainder[3])
if err != nil {
return platform.UDPBindOptions{}, err
}
device, err := decodeBindDevice(remainder[5:])
if err != nil {
return platform.UDPBindOptions{}, err
}
return platform.UDPBindOptions{
Context: socketContext,
LocalAddr: local,
BindDevice: device,
ReuseAddr: reuseAddr,
ReusePort: reusePort,
OnlyV6: onlyV6,
Purpose: purpose,
}, nil
}
func decodeTCPListenOptions(encoded []byte) (platform.TCPListenOptions, error) {
if len(encoded) < 49 || encoded[0] != optionsWireVersion {
return platform.TCPListenOptions{}, fmt.Errorf("invalid TCP listen options")
}
local, err := decodeSocketAddress(encoded[1:28], false)
if err != nil || local == nil {
return platform.TCPListenOptions{}, fmt.Errorf("invalid TCP listen address")
}
socketContext, remainder, err := decodeSocketContext(encoded[28:])
if err != nil {
return platform.TCPListenOptions{}, fmt.Errorf("invalid TCP listen context: %w", err)
}
if len(remainder) < 10 {
return platform.TCPListenOptions{}, fmt.Errorf("truncated TCP listen bind policy")
}
bind, err := decodeTCPBindPolicy(
&net.TCPAddr{IP: local.IP, Port: local.Port, Zone: local.Zone},
socketContext,
remainder[0],
remainder[1],
remainder[2],
remainder[5:],
)
if err != nil {
return platform.TCPListenOptions{}, err
}
purpose, err := decodeTCPListenPurpose(remainder[3])
if err != nil {
return platform.TCPListenOptions{}, err
}
return platform.TCPListenOptions{Bind: bind, Purpose: purpose}, nil
}
func decodeTCPBindPolicy(
localAddr *net.TCPAddr,
socketContext platform.SocketContext,
reuseMode byte,
reusePortByte byte,
onlyV6Byte byte,
deviceBytes []byte,
) (platform.TCPBindOptions, error) {
var reuseAddr *bool
switch reuseMode {
case 0:
case 1:
value := false
reuseAddr = &value
case 2:
value := true
reuseAddr = &value
default:
return platform.TCPBindOptions{}, fmt.Errorf("invalid TCP reuse_addr")
}
reusePort, err := decodeWireBool("TCP reuse_port", reusePortByte)
if err != nil {
return platform.TCPBindOptions{}, err
}
onlyV6, err := decodeWireBool("TCP only_v6", onlyV6Byte)
if err != nil {
return platform.TCPBindOptions{}, err
}
device, err := decodeBindDevice(deviceBytes)
if err != nil {
return platform.TCPBindOptions{}, err
}
return platform.TCPBindOptions{
Context: socketContext,
LocalAddr: localAddr,
BindDevice: device,
ReuseAddr: reuseAddr,
ReusePort: reusePort,
OnlyV6: onlyV6,
}, nil
}
func decodeSocketContext(
encoded []byte,
) (platform.SocketContext, []byte, error) {
if len(encoded) < 11 {
return platform.SocketContext{}, nil, fmt.Errorf("truncated socket context")
}
ipVersion, err := decodeIPVersion(encoded[0])
if err != nil {
return platform.SocketContext{}, nil, err
}
mark, err := decodeSocketMark(encoded[1], encoded[2:6])
if err != nil {
return platform.SocketContext{}, nil, err
}
if encoded[6] > 1 {
return platform.SocketContext{}, nil, fmt.Errorf("invalid netns presence")
}
length := int(binary.BigEndian.Uint32(encoded[7:11]))
if length > len(encoded)-11 || (encoded[6] == 0 && length != 0) {
return platform.SocketContext{}, nil, fmt.Errorf("invalid netns length")
}
var netns *string
if encoded[6] == 1 {
token := encoded[11 : 11+length]
if !utf8.Valid(token) {
return platform.SocketContext{}, nil, fmt.Errorf("netns token is not UTF-8")
}
value := string(token)
netns = &value
}
return platform.SocketContext{
IPVersion: ipVersion,
SocketMark: mark,
NetNS: netns,
}, encoded[11+length:], nil
}
func decodeIPVersion(encoded byte) (platform.IPVersion, error) {
switch encoded {
case 0:
return platform.IPVersionV4, nil
case 1:
return platform.IPVersionV6, nil
case 2:
return platform.IPVersionBoth, nil
default:
return 0, fmt.Errorf("invalid IP version %d", encoded)
}
}
func decodeSocketMark(present byte, encoded []byte) (*uint32, error) {
if present > 1 || len(encoded) != 4 ||
(present == 0 && binary.BigEndian.Uint32(encoded) != 0) {
return nil, fmt.Errorf("invalid socket mark encoding")
}
if present == 0 {
return nil, nil
}
mark := binary.BigEndian.Uint32(encoded)
return &mark, nil
}
func decodeWireBool(name string, encoded byte) (bool, error) {
if encoded > 1 {
return false, fmt.Errorf("invalid %s", name)
}
return encoded == 1, nil
}
func decodeBindDevice(encoded []byte) (*string, error) {
if len(encoded) < 5 || encoded[0] > 1 {
return nil, fmt.Errorf("invalid bind device encoding")
}
length := int(binary.BigEndian.Uint32(encoded[1:5]))
if len(encoded) != 5+length || (encoded[0] == 0 && length != 0) {
return nil, fmt.Errorf("invalid bind device length")
}
if encoded[0] == 0 {
return nil, nil
}
device := string(encoded[5:])
return &device, nil
}
func decodeSocketAddress(encoded []byte, optional bool) (*net.UDPAddr, error) {
if len(encoded) != socketAddressLen {
return nil, fmt.Errorf("invalid socket address length")
}
if optional && encoded[0] == 0 {
for _, value := range encoded[1:] {
if value != 0 {
return nil, fmt.Errorf("noncanonical absent socket address")
}
}
return nil, nil
}
metadata := make([]byte, udpMetadataLen)
copy(metadata, encoded)
address, _, flowinfo, _, err := decodeUDPMetadata(metadata)
if err == nil && flowinfo != 0 {
return nil, fmt.Errorf("IPv6 flowinfo is not supported")
}
return address, err
}
func udpToTCPAddr(address *net.UDPAddr) *net.TCPAddr {
if address == nil {
return nil
}
return &net.TCPAddr{IP: address.IP, Port: address.Port, Zone: address.Zone}
}
func encodeTCPSocketResult(
handle uint64,
localAddr net.Addr,
peerAddr net.Addr,
) ([tcpSocketResultLen]byte, error) {
var encoded [tcpSocketResultLen]byte
binary.BigEndian.PutUint64(encoded[:8], handle)
local, err := encodeNetAddr(localAddr)
if err != nil {
return encoded, err
}
peer, err := encodeNetAddr(peerAddr)
if err != nil {
return encoded, err
}
copy(encoded[8:35], local[:])
copy(encoded[35:], peer[:])
return encoded, nil
}
func encodeBoundSocketResult(
handle uint64,
localAddr net.Addr,
) ([boundSocketResultLen]byte, error) {
var encoded [boundSocketResultLen]byte
binary.BigEndian.PutUint64(encoded[:8], handle)
local, err := encodeNetAddr(localAddr)
if err != nil {
return encoded, err
}
copy(encoded[8:], local[:])
return encoded, nil
}
func encodeNetAddr(address net.Addr) ([socketAddressLen]byte, error) {
var encoded [socketAddressLen]byte
var udpAddr *net.UDPAddr
switch address := address.(type) {
case *net.TCPAddr:
udpAddr = &net.UDPAddr{IP: address.IP, Port: address.Port, Zone: address.Zone}
case *net.UDPAddr:
udpAddr = address
default:
return encoded, fmt.Errorf("unsupported socket address %T", address)
}
metadata, err := encodeUDPMetadata(udpAddr, nil, 0)
if err != nil {
return encoded, err
}
copy(encoded[:], metadata[:socketAddressLen])
return encoded, nil
}
func decodeTCPConnectPurpose(encoded byte) (platform.TCPConnectPurpose, error) {
switch encoded {
case 0:
return platform.TCPConnectDirect, nil
case 1:
return platform.TCPConnectFake, nil
case 2:
return platform.TCPConnectHolePunch, nil
case 3:
return platform.TCPConnectManual, nil
case 4:
return platform.TCPConnectProxyNAT, nil
case 5:
return platform.TCPConnectSTUNProbe, nil
case 6:
return platform.TCPConnectSocks5, nil
case 7:
return platform.TCPConnectPortForward, nil
case 8:
return platform.TCPConnectDataPlane, nil
default:
return 0, fmt.Errorf("invalid TCP connect purpose %d", encoded)
}
}
func decodeUDPBindPurpose(encoded byte) (platform.UDPBindPurpose, error) {
switch encoded {
case 0:
return platform.UDPBindHolePunchControl, nil
case 1:
return platform.UDPBindHolePunchCandidate, nil
case 2:
return platform.UDPBindDirect, nil
case 3:
return platform.UDPBindPortBoundListener, nil
case 4:
return platform.UDPBindProxyNAT, nil
case 5:
return platform.UDPBindSTUNProbe, nil
case 6:
return platform.UDPBindSocks5, nil
case 7:
return platform.UDPBindPortForward, nil
case 8:
return platform.UDPBindPortLease, nil
default:
return 0, fmt.Errorf("invalid UDP bind purpose %d", encoded)
}
}
func decodeTCPListenPurpose(encoded byte) (platform.TCPListenPurpose, error) {
switch encoded {
case 0:
return platform.TCPListenDirect, nil
case 1:
return platform.TCPListenHolePunch, nil
case 2:
return platform.TCPListenManual, nil
case 3:
return platform.TCPListenProxyNAT, nil
case 4:
return platform.TCPListenSocks5, nil
case 5:
return platform.TCPListenPortForward, nil
case 6:
return platform.TCPListenPortLease, nil
default:
return 0, fmt.Errorf("invalid TCP listen purpose %d", encoded)
}
}