188 lines
5.4 KiB
Go
188 lines
5.4 KiB
Go
package e2e
|
|
|
|
import (
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/iceBear67/redapricot/client"
|
|
"github.com/iceBear67/redapricot/client/wire"
|
|
)
|
|
|
|
// blackholeRelay forwards TCP to the hub until it is switched off, after which
|
|
// already-established pairs silently stop carrying bytes while their sockets
|
|
// stay open. That is what a stateful middlebox looks like when it forgets a
|
|
// flow: conntrack expiry, a firewall reload, or a cloud LB idle timeout. No FIN
|
|
// and no RST ever reach either end, so nothing below the application layer can
|
|
// notice.
|
|
// Only flows that already existed when the switch is thrown go dark; new
|
|
// connections are carried normally, exactly as when a middlebox forgets
|
|
// established state but keeps forwarding fresh traffic.
|
|
type blackholeRelay struct {
|
|
addr string
|
|
backend string
|
|
gen atomic.Uint64
|
|
down atomic.Bool
|
|
mu sync.Mutex
|
|
held []net.Conn
|
|
}
|
|
|
|
func newBlackholeRelay(t *testing.T, backend string) *blackholeRelay {
|
|
t.Helper()
|
|
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatalf("relay listen: %v", err)
|
|
}
|
|
r := &blackholeRelay{addr: ln.Addr().String(), backend: backend}
|
|
t.Cleanup(func() {
|
|
_ = ln.Close()
|
|
r.mu.Lock()
|
|
for _, c := range r.held {
|
|
_ = c.Close()
|
|
}
|
|
r.mu.Unlock()
|
|
})
|
|
go func() {
|
|
for {
|
|
cli, err := ln.Accept()
|
|
if err != nil {
|
|
return
|
|
}
|
|
go r.handle(cli)
|
|
}
|
|
}()
|
|
return r
|
|
}
|
|
|
|
func (r *blackholeRelay) handle(cli net.Conn) {
|
|
born := r.gen.Load()
|
|
if r.down.Load() {
|
|
_ = cli.Close()
|
|
return
|
|
}
|
|
up, err := net.Dial("tcp", r.backend)
|
|
if err != nil {
|
|
_ = cli.Close()
|
|
return
|
|
}
|
|
r.mu.Lock()
|
|
r.held = append(r.held, cli, up)
|
|
r.mu.Unlock()
|
|
pipe := func(dst, src net.Conn) {
|
|
buf := make([]byte, 32*1024)
|
|
for {
|
|
n, err := src.Read(buf)
|
|
if n > 0 && r.gen.Load() == born {
|
|
if _, werr := dst.Write(buf[:n]); werr != nil {
|
|
return
|
|
}
|
|
}
|
|
if err != nil {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
go pipe(up, cli)
|
|
go pipe(cli, up)
|
|
}
|
|
|
|
// blackhole strands every currently-established pair. Later connections are
|
|
// unaffected.
|
|
func (r *blackholeRelay) blackhole() { r.gen.Add(1) }
|
|
|
|
// dropAll hard-resets every currently-established pair: a real FIN/RST reaches
|
|
// both ends immediately, as when a middlebox is restarted or a route flaps,
|
|
// rather than the silent stranding blackhole models. Later connections are
|
|
// carried normally.
|
|
//
|
|
// stop takes the relay out of service: new connections are refused rather than
|
|
// carried, so the tunnel stays down until restore is called. Models an outage
|
|
// the client cannot immediately reconnect through.
|
|
func (r *blackholeRelay) stop() { r.down.Store(true) }
|
|
|
|
// restore puts the relay back in service. Flows stranded before it are still
|
|
// dead — only fresh connections are carried, which is what a client gets after
|
|
// a middlebox or upstream link comes back.
|
|
func (r *blackholeRelay) restore() { r.down.Store(false) }
|
|
|
|
// This is the fast path into stream resumption: the client learns the conn is
|
|
// gone at once instead of waiting out a heartbeat timeout.
|
|
func (r *blackholeRelay) dropAll() {
|
|
r.gen.Add(1)
|
|
r.mu.Lock()
|
|
held := r.held
|
|
r.held = nil
|
|
r.mu.Unlock()
|
|
for _, c := range held {
|
|
_ = c.Close()
|
|
}
|
|
}
|
|
|
|
// TestBlackholedPathRecovers is the end-to-end regression guard for the
|
|
// stability bug this hardening was written for: with the tunnel's path silently
|
|
// dropped, the client used to notice nothing at all. Its read loops parked
|
|
// forever, the dead worker conn stayed in the pool, the hub kept routing players
|
|
// to a control session nobody was reading, and no player could connect again
|
|
// until the client process was restarted.
|
|
//
|
|
// The heartbeats must now detect the silence, drop both sessions, and let the
|
|
// existing reconnect path restore service on its own.
|
|
func TestBlackholedPathRecovers(t *testing.T) {
|
|
const psk = "e2e-blackhole"
|
|
hubPort := freePort(t)
|
|
hubAddr := fmt.Sprintf("127.0.0.1:%d", hubPort)
|
|
startHub(t, hubPort, psk)
|
|
dest := newMockDest(t, modeEcho)
|
|
|
|
// The client reaches the hub only through the relay; players connect to the
|
|
// hub directly, as they would from the internet.
|
|
relay := newBlackholeRelay(t, hubAddr)
|
|
const pingMs = 400 // heartbeat timeout is 3x this
|
|
c := startClientWithPing(t, relay.addr, psk, 4, pingMs, []client.Mapping{
|
|
{Pattern: "mc.local", Destination: dest.addr},
|
|
})
|
|
|
|
play := func(what string) error {
|
|
pc, err := net.DialTimeout("tcp", hubAddr, 5*time.Second)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer pc.Close()
|
|
if _, err := pc.Write(wire.BuildHandshake(767, "mc.local", 25565, 2)); err != nil {
|
|
return err
|
|
}
|
|
msg := []byte(what)
|
|
if _, err := pc.Write(msg); err != nil {
|
|
return err
|
|
}
|
|
got := make([]byte, len(msg))
|
|
_ = pc.SetReadDeadline(time.Now().Add(15 * time.Second))
|
|
_, err = io.ReadFull(pc, got)
|
|
return err
|
|
}
|
|
|
|
if err := play("before"); err != nil {
|
|
t.Fatalf("baseline round-trip failed: %v", err)
|
|
}
|
|
|
|
relay.blackhole()
|
|
t.Log("path blackholed: no FIN, no RST, sockets held open")
|
|
|
|
// Wait for the heartbeats to fire, the sessions to be dropped, and the
|
|
// control session to reconnect through a fresh relay pair.
|
|
deadline := time.Now().Add(45 * time.Second)
|
|
var lastErr error
|
|
for time.Now().Before(deadline) {
|
|
if lastErr = play("after"); lastErr == nil {
|
|
t.Logf("recovered on its own; worker conns now %d", c.WorkerConnCount())
|
|
return
|
|
}
|
|
time.Sleep(500 * time.Millisecond)
|
|
}
|
|
t.Fatalf("client never recovered from the blackholed path (last error: %v)", lastErr)
|
|
}
|