less probe
This commit is contained in:
parent
43fe9cf649
commit
9acec15a0a
3 changed files with 351 additions and 74 deletions
41
README.md
41
README.md
|
|
@ -59,8 +59,10 @@ go build -o /usr/local/bin/tsproxy .
|
||||||
| `--name` | `-n` | `tsproxy` | Hostname advertised on the tailnet. |
|
| `--name` | `-n` | `tsproxy` | Hostname advertised on the tailnet. |
|
||||||
| `--dir` | | `~/.config/tsproxy/<name>` | State directory (node identity, keys). |
|
| `--dir` | | `~/.config/tsproxy/<name>` | State directory (node identity, keys). |
|
||||||
| `--verbose` | `-v` | `false` | Verbose tsnet logging. |
|
| `--verbose` | `-v` | `false` | Verbose tsnet logging. |
|
||||||
| `--probe-interval` | | `5s` | How often each target is probed for reachability. |
|
| `--probe-interval` | | `5s` | How often a *down* target is re-probed. No probing happens while it is up. |
|
||||||
| `--probe-timeout` | | `3s` | How long a probe may take before the target counts as unreachable. |
|
| `--probe-timeout` | | `3s` | How long a probe or a forwarded dial may take before it counts as failed. |
|
||||||
|
| `--idle-timeout` | | `5m` | Close a forwarded connection after this long with no traffic either way (`0` disables). |
|
||||||
|
| `--trace` | | `false` | Log every accept, dial, probe and teardown, for diagnosing stalls. |
|
||||||
|
|
||||||
Target can be any MagicDNS name, short hostname, or tailnet IP.
|
Target can be any MagicDNS name, short hostname, or tailnet IP.
|
||||||
|
|
||||||
|
|
@ -93,11 +95,18 @@ plist, then remove it once the state directory has been populated.
|
||||||
- **Startup is strict.** Every local address is bound once at startup as a
|
- **Startup is strict.** Every local address is bound once at startup as a
|
||||||
check; if any fails, the process exits — partial success is confusing under
|
check; if any fails, the process exits — partial success is confusing under
|
||||||
a supervisor.
|
a supervisor.
|
||||||
- **The local port tracks the target.** Each forward dials its target every
|
- **The local port tracks the target.** A forward only keeps its local port
|
||||||
`--probe-interval` and only keeps the local port bound while that succeeds.
|
bound while the target is believed reachable. So an unavailable target means
|
||||||
So an unavailable target means `connection refused` on the local port, not a
|
`connection refused` on the local port, not a connect that immediately EOFs,
|
||||||
connect that immediately EOFs, and clients back off the way they would
|
and clients back off the way they would against a genuinely down service.
|
||||||
against a genuinely down service.
|
- **Probing only happens while down.** One probe runs at startup to decide
|
||||||
|
whether to bind at all. After that, real connections are the health signal
|
||||||
|
and nothing polls — a healthy target is never dialled except to carry
|
||||||
|
traffic, so services that log every connection stay quiet. Probing resumes
|
||||||
|
(every `--probe-interval`) only once something has failed, and stops again as
|
||||||
|
soon as the target answers. The cost of not polling is that a target which
|
||||||
|
dies unnoticed leaves the port bound until something tries to use it: that
|
||||||
|
one connection is accepted and then reset, and the port closes behind it.
|
||||||
- **One failed connection closes the port.** A client that reconnects the
|
- **One failed connection closes the port.** A client that reconnects the
|
||||||
instant its connection breaks would otherwise beat the next probe and be
|
instant its connection breaks would otherwise beat the next probe and be
|
||||||
accepted into a forward with nothing behind it. So any connection that fails
|
accepted into a forward with nothing behind it. So any connection that fails
|
||||||
|
|
@ -113,9 +122,21 @@ plist, then remove it once the state directory has been populated.
|
||||||
probe succeed. Forwarded connections use the same timeout as probes.
|
probe succeed. Forwarded connections use the same timeout as probes.
|
||||||
- **Target loss drops live connections.** A tailnet peer can vanish without the
|
- **Target loss drops live connections.** A tailnet peer can vanish without the
|
||||||
userspace TCP stack ever erroring on an established connection, which leaves
|
userspace TCP stack ever erroring on an established connection, which leaves
|
||||||
local sockets hanging and apps waiting on a dead link. When a probe fails,
|
local sockets hanging and apps waiting on a dead link. Two things catch this.
|
||||||
every connection on that forward is closed so clients see the drop and
|
A connection that fails outright takes the port down immediately, and the
|
||||||
reconnect. Detection takes up to `--probe-interval` + `--probe-timeout`.
|
probe that follows resets everything still in flight. A connection that
|
||||||
|
merely goes silent is caught by `--idle-timeout`, since with polling switched
|
||||||
|
off there is nothing else watching it.
|
||||||
|
- **Idle connections are closed.** Once `--idle-timeout` passes with no bytes
|
||||||
|
moving *either way*, the connection is reset. Traffic in one direction counts,
|
||||||
|
so a long upload or download is never reaped mid-stream. This is what bounds
|
||||||
|
the case where the target accepts a connection and then goes silent — but it
|
||||||
|
cannot tell that apart from a connection legitimately sitting idle, so a
|
||||||
|
shell session or a database pool left quiet past the timeout is dropped too
|
||||||
|
and has to reconnect. Lower it to notice dead targets sooner, raise it if
|
||||||
|
your clients hold connections open across long gaps, `0` to disable. Reaping
|
||||||
|
an idle connection never closes the local port: an unused connection says
|
||||||
|
nothing about whether the target is healthy.
|
||||||
- **Failures are reset, not closed.** How a connection ends is forwarded
|
- **Failures are reset, not closed.** How a connection ends is forwarded
|
||||||
faithfully. A target that closes cleanly gives the local client a FIN (an
|
faithfully. A target that closes cleanly gives the local client a FIN (an
|
||||||
ordinary EOF); a target that resets, errors, or disappears gives it an RST.
|
ordinary EOF); a target that resets, errors, or disappears gives it an RST.
|
||||||
|
|
|
||||||
226
main.go
226
main.go
|
|
@ -8,6 +8,7 @@ package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"log"
|
"log"
|
||||||
"net"
|
"net"
|
||||||
|
|
@ -15,6 +16,7 @@ import (
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
flag "github.com/spf13/pflag"
|
flag "github.com/spf13/pflag"
|
||||||
|
|
@ -52,9 +54,13 @@ func main() {
|
||||||
dir = flag.String("dir", "", "state directory (default: ~/.config/tsproxy/<name>)")
|
dir = flag.String("dir", "", "state directory (default: ~/.config/tsproxy/<name>)")
|
||||||
verbose = flag.BoolP("verbose", "v", false, "verbose tsnet logging")
|
verbose = flag.BoolP("verbose", "v", false, "verbose tsnet logging")
|
||||||
interval = flag.Duration("probe-interval", 5*time.Second,
|
interval = flag.Duration("probe-interval", 5*time.Second,
|
||||||
"how often to probe each target for reachability")
|
"how often to re-probe a target that is down (a target that is up is never probed)")
|
||||||
timeout = flag.Duration("probe-timeout", 3*time.Second,
|
timeout = flag.Duration("probe-timeout", 3*time.Second,
|
||||||
"how long a target probe may take before the target counts as unreachable")
|
"how long a probe or a forwarded dial may take before it counts as failed")
|
||||||
|
trace = flag.Bool("trace", false,
|
||||||
|
"log every accept, dial, probe and teardown (for diagnosing stalls)")
|
||||||
|
idle = flag.Duration("idle-timeout", 5*time.Minute,
|
||||||
|
"close a forwarded connection after this long with no traffic either way (0 disables)")
|
||||||
)
|
)
|
||||||
flag.Parse()
|
flag.Parse()
|
||||||
|
|
||||||
|
|
@ -71,6 +77,9 @@ func main() {
|
||||||
if *timeout <= 0 {
|
if *timeout <= 0 {
|
||||||
log.Fatal("--probe-timeout must be positive")
|
log.Fatal("--probe-timeout must be positive")
|
||||||
}
|
}
|
||||||
|
if *idle < 0 {
|
||||||
|
log.Fatal("--idle-timeout must not be negative (0 disables)")
|
||||||
|
}
|
||||||
|
|
||||||
stateDir := *dir
|
stateDir := *dir
|
||||||
if stateDir == "" {
|
if stateDir == "" {
|
||||||
|
|
@ -120,6 +129,8 @@ func main() {
|
||||||
timeout: *timeout,
|
timeout: *timeout,
|
||||||
recheck: make(chan struct{}, 1),
|
recheck: make(chan struct{}, 1),
|
||||||
conns: make(map[net.Conn]struct{}),
|
conns: make(map[net.Conn]struct{}),
|
||||||
|
trace: *trace,
|
||||||
|
idle: *idle,
|
||||||
}
|
}
|
||||||
log.Printf("tsproxy: %s -> %s (via tailnet as %q)", f.local, f.target, *hostname)
|
log.Printf("tsproxy: %s -> %s (via tailnet as %q)", f.local, f.target, *hostname)
|
||||||
wg.Add(1)
|
wg.Add(1)
|
||||||
|
|
@ -154,7 +165,10 @@ type proxy struct {
|
||||||
target string
|
target string
|
||||||
interval time.Duration
|
interval time.Duration
|
||||||
timeout time.Duration
|
timeout time.Duration
|
||||||
|
idle time.Duration // tear down a connection after this long with no traffic; 0 disables
|
||||||
recheck chan struct{} // nudges the probe loop to re-probe immediately
|
recheck chan struct{} // nudges the probe loop to re-probe immediately
|
||||||
|
trace bool // log every accept, dial, probe and teardown
|
||||||
|
seq atomic.Uint64 // connection counter, for correlating log lines
|
||||||
|
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
state state
|
state state
|
||||||
|
|
@ -162,28 +176,76 @@ type proxy struct {
|
||||||
conns map[net.Conn]struct{}
|
conns map[net.Conn]struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (p *proxy) logf(format string, args ...any) {
|
||||||
|
log.Printf("tsproxy: %s -> %s: %s", p.local, p.target, fmt.Sprintf(format, args...))
|
||||||
|
}
|
||||||
|
|
||||||
|
// tracef logs only under --trace: the per-connection and per-probe detail you
|
||||||
|
// want while diagnosing a stall, and not otherwise.
|
||||||
|
func (p *proxy) tracef(format string, args ...any) {
|
||||||
|
if p.trace {
|
||||||
|
p.logf(format, args...)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *proxy) liveConns() int {
|
||||||
|
p.mu.Lock()
|
||||||
|
defer p.mu.Unlock()
|
||||||
|
return len(p.conns)
|
||||||
|
}
|
||||||
|
|
||||||
|
// run probes the target and binds or unbinds the local port to match.
|
||||||
|
//
|
||||||
|
// Probing happens only while the target is believed down: once it is up, real
|
||||||
|
// connections are the health signal and further probes would be pure noise
|
||||||
|
// against the target. One probe runs at startup to decide whether to bind at
|
||||||
|
// all, and after that the loop sits idle until something reports a failure.
|
||||||
func (p *proxy) run(ctx context.Context) {
|
func (p *proxy) run(ctx context.Context) {
|
||||||
t := time.NewTicker(p.interval)
|
|
||||||
defer t.Stop()
|
|
||||||
for {
|
for {
|
||||||
if err := p.probe(ctx); err != nil {
|
start := time.Now()
|
||||||
|
err := p.probe(ctx)
|
||||||
|
took := time.Since(start).Round(time.Millisecond)
|
||||||
|
if err != nil {
|
||||||
|
p.tracef("probe failed after %s: %v (%d live)", took, err, p.liveConns())
|
||||||
p.markDown(err)
|
p.markDown(err)
|
||||||
} else {
|
} else {
|
||||||
|
p.tracef("probe ok in %s (%d live)", took, p.liveConns())
|
||||||
p.markUp(ctx)
|
p.markUp(ctx)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Retry on a timer only while down. While up, wait to be woken by a
|
||||||
|
// connection that failed -- there is nothing to poll for.
|
||||||
|
var retry <-chan time.Time
|
||||||
|
var timer *time.Timer
|
||||||
|
if !p.isUp() {
|
||||||
|
timer = time.NewTimer(p.interval)
|
||||||
|
retry = timer.C
|
||||||
|
}
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
|
if timer != nil {
|
||||||
|
timer.Stop()
|
||||||
|
}
|
||||||
p.markDown(ctx.Err())
|
p.markDown(ctx.Err())
|
||||||
return
|
return
|
||||||
case <-t.C:
|
case <-retry:
|
||||||
case <-p.recheck:
|
case <-p.recheck:
|
||||||
|
if timer != nil {
|
||||||
|
timer.Stop()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// probe dials the target and hangs up. It is the only evidence we have that
|
func (p *proxy) isUp() bool {
|
||||||
// the target is alive: a tailnet peer can disappear without the userspace TCP
|
p.mu.Lock()
|
||||||
// stack ever reporting an error on an established connection.
|
defer p.mu.Unlock()
|
||||||
|
return p.state == stateUp
|
||||||
|
}
|
||||||
|
|
||||||
|
// probe dials the target and hangs up. It is how a down target is found to be
|
||||||
|
// back: a tailnet peer can disappear without the userspace TCP stack ever
|
||||||
|
// reporting an error on an established connection.
|
||||||
func (p *proxy) probe(ctx context.Context) error {
|
func (p *proxy) probe(ctx context.Context) error {
|
||||||
ctx, cancel := context.WithTimeout(ctx, p.timeout)
|
ctx, cancel := context.WithTimeout(ctx, p.timeout)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
@ -224,9 +286,9 @@ func (p *proxy) markUp(ctx context.Context) {
|
||||||
p.ln = ln
|
p.ln = ln
|
||||||
if was == stateUp {
|
if was == stateUp {
|
||||||
// Listener was lost without the target going down; already logged.
|
// Listener was lost without the target going down; already logged.
|
||||||
log.Printf("tsproxy: %s -> %s: listening again", p.local, p.target)
|
p.logf("PORT OPEN: listening again on %s", p.local)
|
||||||
} else {
|
} else {
|
||||||
log.Printf("tsproxy: %s -> %s: target reachable, accepting connections", p.local, p.target)
|
p.logf("PORT OPEN: target reachable, listening on %s", p.local)
|
||||||
}
|
}
|
||||||
go p.accept(ctx, ln)
|
go p.accept(ctx, ln)
|
||||||
}
|
}
|
||||||
|
|
@ -248,8 +310,9 @@ func (p *proxy) markDown(cause error) {
|
||||||
p.mu.Unlock()
|
p.mu.Unlock()
|
||||||
|
|
||||||
if was != stateDown {
|
if was != stateDown {
|
||||||
log.Printf("tsproxy: %s -> %s: target unreachable: %v (refusing connections, reset %d in flight)",
|
p.logf("PORT CLOSED: target unreachable: %v (refusing connections, resetting %d in flight)", cause, len(conns))
|
||||||
p.local, p.target, cause, len(conns))
|
} else {
|
||||||
|
p.tracef("still down: %v (%d in flight to reset)", cause, len(conns))
|
||||||
}
|
}
|
||||||
if ln != nil {
|
if ln != nil {
|
||||||
ln.Close()
|
ln.Close()
|
||||||
|
|
@ -282,8 +345,9 @@ func (p *proxy) suspend() {
|
||||||
ln.Close()
|
ln.Close()
|
||||||
}
|
}
|
||||||
if was == stateUp {
|
if was == stateUp {
|
||||||
log.Printf("tsproxy: %s -> %s: connection to target failed, refusing connections pending probe",
|
p.logf("PORT CLOSED: a connection to the target failed, refusing connections pending probe")
|
||||||
p.local, p.target)
|
} else {
|
||||||
|
p.tracef("suspend: port was already closed")
|
||||||
}
|
}
|
||||||
p.nudge()
|
p.nudge()
|
||||||
}
|
}
|
||||||
|
|
@ -302,19 +366,28 @@ func (p *proxy) accept(ctx context.Context, ln net.Listener) {
|
||||||
}
|
}
|
||||||
p.mu.Unlock()
|
p.mu.Unlock()
|
||||||
if current {
|
if current {
|
||||||
log.Printf("tsproxy: %s -> %s: accept: %v", p.local, p.target, err)
|
// Nobody is accepting on a bound port now -- that would hang a
|
||||||
|
// client in connect(), so say so loudly.
|
||||||
|
p.logf("PORT CLOSED: accept loop died on %s: %v", p.local, err)
|
||||||
ln.Close()
|
ln.Close()
|
||||||
p.nudge()
|
p.nudge()
|
||||||
|
} else {
|
||||||
|
p.tracef("accept loop exited on %s (listener already replaced)", p.local)
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
id := p.seq.Add(1)
|
||||||
|
p.tracef("[conn %d] accepted from %s", id, c.RemoteAddr())
|
||||||
if !p.track(c) {
|
if !p.track(c) {
|
||||||
// Raced with markDown: the target went away between Accept and
|
// Raced with the port closing: the target went down between Accept
|
||||||
// here, so this connection is already condemned.
|
// and here, so this connection is already condemned. Reset rather
|
||||||
c.Close()
|
// than close, or the client sees a successful connect followed by a
|
||||||
|
// clean EOF -- the exact thing the port closing exists to prevent.
|
||||||
|
p.tracef("[conn %d] reset: target went down during accept", id)
|
||||||
|
abort(c)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
go p.handle(ctx, c)
|
go p.handle(ctx, c, id)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -395,41 +468,132 @@ func (t *targetConn) failure() error {
|
||||||
return t.err
|
return t.err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (p *proxy) handle(ctx context.Context, in net.Conn) {
|
// which side of a forwarded connection finished first, and what it moved.
|
||||||
|
type direction struct {
|
||||||
|
name string
|
||||||
|
bytes int64
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
// countingReader records that bytes moved, for the idle watchdog. Both
|
||||||
|
// directions share one counter: traffic either way means the connection is
|
||||||
|
// alive, so an active download must not let the quiet upload side time out.
|
||||||
|
type countingReader struct {
|
||||||
|
r io.Reader
|
||||||
|
moved *atomic.Int64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *countingReader) Read(b []byte) (int, error) {
|
||||||
|
n, err := c.r.Read(b)
|
||||||
|
if n > 0 {
|
||||||
|
c.moved.Add(int64(n))
|
||||||
|
}
|
||||||
|
return n, err
|
||||||
|
}
|
||||||
|
|
||||||
|
// watchIdle tears a connection down once no bytes have moved either way for
|
||||||
|
// p.idle. With probing suppressed while the target is up, this is what catches
|
||||||
|
// a target that accepted a connection and then went silent -- otherwise the
|
||||||
|
// client waits forever on a socket whose far end is gone.
|
||||||
|
//
|
||||||
|
// It sets idled before closing so the teardown can tell this apart from a
|
||||||
|
// target failure: a connection nobody was using is no evidence the target is
|
||||||
|
// down, and must not take the local port with it.
|
||||||
|
func (p *proxy) watchIdle(id uint64, in, out net.Conn, moved *atomic.Int64, idled *atomic.Bool, stop <-chan struct{}) {
|
||||||
|
tick := p.idle / 4
|
||||||
|
if tick <= 0 {
|
||||||
|
tick = p.idle
|
||||||
|
}
|
||||||
|
t := time.NewTicker(tick)
|
||||||
|
defer t.Stop()
|
||||||
|
|
||||||
|
last := moved.Load()
|
||||||
|
lastChange := time.Now()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-stop:
|
||||||
|
return
|
||||||
|
case <-t.C:
|
||||||
|
if n := moved.Load(); n != last {
|
||||||
|
last, lastChange = n, time.Now()
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if quiet := time.Since(lastChange); quiet >= p.idle {
|
||||||
|
idled.Store(true)
|
||||||
|
p.logf("[conn %d] no traffic for %s, closing", id, quiet.Round(time.Second))
|
||||||
|
abort(in)
|
||||||
|
out.Close()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *proxy) handle(ctx context.Context, in net.Conn, id uint64) {
|
||||||
defer p.untrack(in)
|
defer p.untrack(in)
|
||||||
|
opened := time.Now()
|
||||||
|
|
||||||
// Bound the dial. A tailnet peer that is routable but dead accepts nothing
|
// Bound the dial. A tailnet peer that is routable but dead accepts nothing
|
||||||
// and refuses nothing, so an unbounded dial parks here forever, holding a
|
// and refuses nothing, so an unbounded dial parks here forever, holding a
|
||||||
// local socket open with no way out: the probe loop only tears down live
|
// local socket open with no way out: the probe loop only tears down live
|
||||||
// connections when a probe fails, and a target that recovers makes the
|
// connections when a probe fails, and a target that recovers makes the
|
||||||
// probe succeed. The connection would hang for good.
|
// probe succeed. The connection would hang for good.
|
||||||
|
p.tracef("[conn %d] dialing %s (timeout %s)", id, p.target, p.timeout)
|
||||||
|
dialStart := time.Now()
|
||||||
dialCtx, cancel := context.WithTimeout(ctx, p.timeout)
|
dialCtx, cancel := context.WithTimeout(ctx, p.timeout)
|
||||||
c, err := p.dial.Dial(dialCtx, "tcp", p.target)
|
c, err := p.dial.Dial(dialCtx, "tcp", p.target)
|
||||||
cancel() // governs the dial only; the returned conn outlives it
|
cancel() // governs the dial only; the returned conn outlives it
|
||||||
|
dialTook := time.Since(dialStart).Round(time.Millisecond)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("tsproxy: %s -> %s: dial: %v", p.local, p.target, err)
|
p.logf("[conn %d] dial failed after %s: %v -- resetting client", id, dialTook, err)
|
||||||
abort(in)
|
abort(in)
|
||||||
p.suspend()
|
p.suspend()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
p.tracef("[conn %d] connected to target in %s", id, dialTook)
|
||||||
out := &targetConn{Conn: c}
|
out := &targetConn{Conn: c}
|
||||||
defer c.Close()
|
defer c.Close()
|
||||||
|
|
||||||
|
var moved atomic.Int64
|
||||||
|
var idled atomic.Bool
|
||||||
|
if p.idle > 0 {
|
||||||
|
stop := make(chan struct{})
|
||||||
|
defer close(stop)
|
||||||
|
go p.watchIdle(id, in, out, &moved, &idled, stop)
|
||||||
|
}
|
||||||
|
|
||||||
// A clean EOF from the target is forwarded as a FIN -- exactly what the
|
// A clean EOF from the target is forwarded as a FIN -- exactly what the
|
||||||
// client would have seen connecting directly. Anything else is forwarded
|
// client would have seen connecting directly. Anything else is forwarded
|
||||||
// as an RST, and if the target was at fault, also takes the port down.
|
// as an RST, and if the target was at fault, also takes the port down.
|
||||||
done := make(chan error, 2)
|
done := make(chan direction, 2)
|
||||||
go func() { _, err := io.Copy(out, in); done <- err }()
|
go func() {
|
||||||
go func() { _, err := io.Copy(in, out); done <- err }()
|
n, err := io.Copy(out, &countingReader{in, &moved})
|
||||||
|
done <- direction{"client->target", n, err}
|
||||||
|
}()
|
||||||
|
go func() {
|
||||||
|
n, err := io.Copy(in, &countingReader{out, &moved})
|
||||||
|
done <- direction{"target->client", n, err}
|
||||||
|
}()
|
||||||
|
|
||||||
copyErr := <-done
|
first := <-done
|
||||||
targetErr := out.failure()
|
targetErr := out.failure()
|
||||||
if copyErr != nil || targetErr != nil {
|
lived := time.Since(opened).Round(time.Millisecond)
|
||||||
|
|
||||||
|
switch {
|
||||||
|
case idled.Load():
|
||||||
|
// The watchdog already reset the client. Deliberately no suspend: an
|
||||||
|
// unused connection says nothing about whether the target is healthy.
|
||||||
|
p.tracef("[conn %d] idle-closed after %s, %d bytes total", id, lived, moved.Load())
|
||||||
|
case targetErr != nil:
|
||||||
|
p.logf("[conn %d] target failed after %s: %v (%s moved %d bytes) -- resetting client",
|
||||||
|
id, lived, targetErr, first.name, first.bytes)
|
||||||
abort(in)
|
abort(in)
|
||||||
} else {
|
p.suspend()
|
||||||
|
case first.err != nil:
|
||||||
|
p.tracef("[conn %d] reset after %s: %s: %v (%d bytes)", id, lived, first.name, first.err, first.bytes)
|
||||||
|
abort(in)
|
||||||
|
default:
|
||||||
|
p.tracef("[conn %d] closed cleanly after %s: %s ended, %d bytes", id, lived, first.name, first.bytes)
|
||||||
in.Close()
|
in.Close()
|
||||||
}
|
}
|
||||||
if targetErr != nil {
|
|
||||||
p.suspend()
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
156
main_test.go
156
main_test.go
|
|
@ -6,6 +6,7 @@ import (
|
||||||
"io"
|
"io"
|
||||||
"net"
|
"net"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"syscall"
|
"syscall"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
@ -26,9 +27,11 @@ const (
|
||||||
// established just sit there with no error from the userspace TCP stack.
|
// established just sit there with no error from the userspace TCP stack.
|
||||||
type fakeTarget struct {
|
type fakeTarget struct {
|
||||||
ln net.Listener
|
ln net.Listener
|
||||||
|
dials atomic.Int64 // every dial, probe or forwarded alike
|
||||||
|
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
mode dialMode
|
mode dialMode
|
||||||
|
sink bool // swallow input and never reply, instead of echoing
|
||||||
accepted []net.Conn
|
accepted []net.Conn
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -48,8 +51,16 @@ func newFakeTarget(t *testing.T) *fakeTarget {
|
||||||
}
|
}
|
||||||
ft.mu.Lock()
|
ft.mu.Lock()
|
||||||
ft.accepted = append(ft.accepted, c)
|
ft.accepted = append(ft.accepted, c)
|
||||||
|
sink := ft.sink
|
||||||
ft.mu.Unlock()
|
ft.mu.Unlock()
|
||||||
go func() { io.Copy(c, c); c.Close() }()
|
go func() {
|
||||||
|
if sink {
|
||||||
|
io.Copy(io.Discard, c) // read forever, never reply
|
||||||
|
} else {
|
||||||
|
io.Copy(c, c)
|
||||||
|
}
|
||||||
|
c.Close()
|
||||||
|
}()
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
return ft
|
return ft
|
||||||
|
|
@ -70,6 +81,13 @@ func (f *fakeTarget) killAccepted(reset bool) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// setSink makes the target read and never reply, so traffic flows one way only.
|
||||||
|
func (f *fakeTarget) setSink(v bool) {
|
||||||
|
f.mu.Lock()
|
||||||
|
defer f.mu.Unlock()
|
||||||
|
f.sink = v
|
||||||
|
}
|
||||||
|
|
||||||
func (f *fakeTarget) setMode(m dialMode) {
|
func (f *fakeTarget) setMode(m dialMode) {
|
||||||
f.mu.Lock()
|
f.mu.Lock()
|
||||||
defer f.mu.Unlock()
|
defer f.mu.Unlock()
|
||||||
|
|
@ -85,6 +103,7 @@ func (f *fakeTarget) setReachable(v bool) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *fakeTarget) Dial(ctx context.Context, network, addr string) (net.Conn, error) {
|
func (f *fakeTarget) Dial(ctx context.Context, network, addr string) (net.Conn, error) {
|
||||||
|
f.dials.Add(1)
|
||||||
f.mu.Lock()
|
f.mu.Lock()
|
||||||
mode := f.mode
|
mode := f.mode
|
||||||
f.mu.Unlock()
|
f.mu.Unlock()
|
||||||
|
|
@ -113,23 +132,7 @@ func freePort(t *testing.T) string {
|
||||||
|
|
||||||
func startProxy(t *testing.T, ft *fakeTarget, local string) *proxy {
|
func startProxy(t *testing.T, ft *fakeTarget, local string) *proxy {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
p := &proxy{
|
return startProxyWith(t, ft, local, 20*time.Millisecond, time.Second, 0)
|
||||||
dial: ft,
|
|
||||||
local: local,
|
|
||||||
target: "target:1234",
|
|
||||||
interval: 20 * time.Millisecond,
|
|
||||||
timeout: time.Second,
|
|
||||||
recheck: make(chan struct{}, 1),
|
|
||||||
conns: make(map[net.Conn]struct{}),
|
|
||||||
}
|
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
|
||||||
done := make(chan struct{})
|
|
||||||
go func() { defer close(done); p.run(ctx) }()
|
|
||||||
t.Cleanup(func() {
|
|
||||||
cancel()
|
|
||||||
<-done
|
|
||||||
})
|
|
||||||
return p
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// waitFor polls cond until it holds or the deadline passes.
|
// waitFor polls cond until it holds or the deadline passes.
|
||||||
|
|
@ -203,10 +206,14 @@ func TestForwardsData(t *testing.T) {
|
||||||
// The reported bug: when the proxy loses the target, an already-established
|
// The reported bug: when the proxy loses the target, an already-established
|
||||||
// local socket must be closed so the client's read fails and it reconnects,
|
// local socket must be closed so the client's read fails and it reconnects,
|
||||||
// rather than blocking forever on a half-dead connection.
|
// rather than blocking forever on a half-dead connection.
|
||||||
func TestTargetLossClosesLocalConnection(t *testing.T) {
|
//
|
||||||
|
// Nothing polls while the target is up, so the idle timeout is what catches
|
||||||
|
// this: the target vanished without any TCP signal, and the connection simply
|
||||||
|
// goes quiet. The teardown must still be a reset, not a clean EOF.
|
||||||
|
func TestSilentTargetLossClosesIdleConnection(t *testing.T) {
|
||||||
ft := newFakeTarget(t)
|
ft := newFakeTarget(t)
|
||||||
local := freePort(t)
|
local := freePort(t)
|
||||||
startProxy(t, ft, local)
|
startProxyWith(t, ft, local, 30*time.Second, time.Second, 250*time.Millisecond)
|
||||||
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
c, err := net.Dial("tcp", local)
|
c, err := net.Dial("tcp", local)
|
||||||
|
|
@ -225,26 +232,109 @@ func TestTargetLossClosesLocalConnection(t *testing.T) {
|
||||||
t.Fatalf("read echo: %v", err)
|
t.Fatalf("read echo: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Target vanishes. The echo server still holds the far end open, so
|
// Target vanishes. The echo server still holds the far end open, so there
|
||||||
// nothing but the probe loop can notice.
|
// is no TCP signal at all -- the connection just stops carrying traffic.
|
||||||
ft.setReachable(false)
|
ft.setReachable(false)
|
||||||
|
|
||||||
// The client's blocking read must return, promptly, and as a reset -- a
|
// The client's blocking read must return, and as a reset: a plain EOF here
|
||||||
// plain EOF here would tell the client the stream ended normally.
|
// would tell the client the stream ended normally.
|
||||||
c.SetReadDeadline(time.Now().Add(3 * time.Second))
|
c.SetReadDeadline(time.Now().Add(5 * time.Second))
|
||||||
_, err = c.Read(buf)
|
_, err = c.Read(buf)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Fatal("read succeeded after target loss; want the connection closed")
|
t.Fatal("read succeeded after target loss; want the connection closed")
|
||||||
}
|
}
|
||||||
if isTimeout(err) {
|
if isTimeout(err) {
|
||||||
t.Fatal("read blocked after target loss; local socket was never closed")
|
t.Fatal("read blocked after target loss; the idle timeout never fired")
|
||||||
}
|
}
|
||||||
if !isReset(err) {
|
if !isReset(err) {
|
||||||
t.Errorf("read err = %v, want a connection reset (a clean EOF would falsely signal a complete stream)", err)
|
t.Errorf("read err = %v, want a connection reset (a clean EOF would falsely signal a complete stream)", err)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// New connects must be refused too.
|
// An idle connection is not evidence that the target is unhealthy, so reaping
|
||||||
waitFor(t, "local port to be refused", func() bool { return localPortRefused(local) })
|
// one must not take the local port down with it.
|
||||||
|
func TestIdleTimeoutLeavesPortOpen(t *testing.T) {
|
||||||
|
ft := newFakeTarget(t)
|
||||||
|
local := freePort(t)
|
||||||
|
startProxyWith(t, ft, local, 30*time.Second, time.Second, 200*time.Millisecond)
|
||||||
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
|
c, err := net.Dial("tcp", local)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial local: %v", err)
|
||||||
|
}
|
||||||
|
defer c.Close()
|
||||||
|
|
||||||
|
// Never send anything; let the watchdog reap it.
|
||||||
|
c.SetReadDeadline(time.Now().Add(5 * time.Second))
|
||||||
|
buf := make([]byte, 1)
|
||||||
|
if _, err := c.Read(buf); err == nil || isTimeout(err) {
|
||||||
|
t.Fatalf("idle connection was not reaped (err: %v)", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The target is perfectly healthy, so the port must still be accepting.
|
||||||
|
// A long probe interval means a wrongly-closed port would stay closed.
|
||||||
|
if localPortRefused(local) {
|
||||||
|
t.Fatal("port closed after an idle reap; an unused connection says nothing about the target")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Traffic in one direction must keep the connection alive even though the other
|
||||||
|
// direction is silent, or every upload and every long download would be reaped
|
||||||
|
// mid-stream. The target here never replies, so client->target is the only
|
||||||
|
// activity there is -- an idle check that watched just one direction would kill
|
||||||
|
// this connection.
|
||||||
|
func TestOneWayTrafficIsNotIdle(t *testing.T) {
|
||||||
|
ft := newFakeTarget(t)
|
||||||
|
ft.setSink(true)
|
||||||
|
local := freePort(t)
|
||||||
|
idle := 200 * time.Millisecond
|
||||||
|
startProxyWith(t, ft, local, 30*time.Second, time.Second, idle)
|
||||||
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
|
c, err := net.Dial("tcp", local)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("dial local: %v", err)
|
||||||
|
}
|
||||||
|
defer c.Close()
|
||||||
|
|
||||||
|
// Write steadily for well over the idle window without ever reading.
|
||||||
|
deadline := time.Now().Add(4 * idle)
|
||||||
|
for time.Now().Before(deadline) {
|
||||||
|
c.SetWriteDeadline(time.Now().Add(time.Second))
|
||||||
|
if _, err := c.Write([]byte("x")); err != nil {
|
||||||
|
t.Fatalf("connection died under one-way traffic within %s: %v", idle, err)
|
||||||
|
}
|
||||||
|
time.Sleep(idle / 8)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Still writable, i.e. still alive after 4x the idle window of one-way use.
|
||||||
|
c.SetWriteDeadline(time.Now().Add(time.Second))
|
||||||
|
if _, err := c.Write([]byte("x")); err != nil {
|
||||||
|
t.Fatalf("connection unusable after sustained one-way traffic: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The point of the change: a healthy target must never be dialled again after
|
||||||
|
// the one probe that decides whether to bind at startup.
|
||||||
|
func TestNoProbingWhileTargetIsUp(t *testing.T) {
|
||||||
|
ft := newFakeTarget(t)
|
||||||
|
local := freePort(t)
|
||||||
|
interval := 20 * time.Millisecond
|
||||||
|
startProxyWith(t, ft, local, interval, time.Second, 0)
|
||||||
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
|
// The readiness check above opens a real connection, which forwards and so
|
||||||
|
// dials the target. Let that finish before snapshotting, or its dial lands
|
||||||
|
// inside the measurement window and looks like a probe.
|
||||||
|
time.Sleep(5 * interval)
|
||||||
|
|
||||||
|
settled := ft.dials.Load()
|
||||||
|
time.Sleep(20 * interval) // twenty probe intervals' worth of opportunity
|
||||||
|
if got := ft.dials.Load(); got != settled {
|
||||||
|
t.Errorf("target was dialled %d more times while up; want 0 (probing must stop once up)", got-settled)
|
||||||
|
}
|
||||||
|
t.Logf("%d dials total while up over %s", ft.dials.Load(), 20*interval)
|
||||||
}
|
}
|
||||||
|
|
||||||
func isTimeout(err error) bool {
|
func isTimeout(err error) bool {
|
||||||
|
|
@ -325,7 +415,7 @@ func TestTargetHardCloseIsNoticedImmediately(t *testing.T) {
|
||||||
|
|
||||||
// startProxyWith runs a proxy with explicit timings, for tests that need the
|
// startProxyWith runs a proxy with explicit timings, for tests that need the
|
||||||
// probe loop held back so a pass can only come from the connection path.
|
// probe loop held back so a pass can only come from the connection path.
|
||||||
func startProxyWith(t *testing.T, ft *fakeTarget, local string, interval, timeout time.Duration) *proxy {
|
func startProxyWith(t *testing.T, ft *fakeTarget, local string, interval, timeout, idle time.Duration) *proxy {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
p := &proxy{
|
p := &proxy{
|
||||||
dial: ft,
|
dial: ft,
|
||||||
|
|
@ -333,8 +423,10 @@ func startProxyWith(t *testing.T, ft *fakeTarget, local string, interval, timeou
|
||||||
target: "target:1234",
|
target: "target:1234",
|
||||||
interval: interval,
|
interval: interval,
|
||||||
timeout: timeout,
|
timeout: timeout,
|
||||||
|
idle: idle,
|
||||||
recheck: make(chan struct{}, 1),
|
recheck: make(chan struct{}, 1),
|
||||||
conns: make(map[net.Conn]struct{}),
|
conns: make(map[net.Conn]struct{}),
|
||||||
|
trace: true, // so a failing test leaves a usable log behind
|
||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
done := make(chan struct{})
|
done := make(chan struct{})
|
||||||
|
|
@ -353,7 +445,7 @@ func TestHangingDialDoesNotStrandClient(t *testing.T) {
|
||||||
local := freePort(t)
|
local := freePort(t)
|
||||||
// Probe interval far beyond the test, so the probe loop cannot be what
|
// Probe interval far beyond the test, so the probe loop cannot be what
|
||||||
// rescues the connection.
|
// rescues the connection.
|
||||||
startProxyWith(t, ft, local, 30*time.Second, 300*time.Millisecond)
|
startProxyWith(t, ft, local, 30*time.Second, 300*time.Millisecond, 0)
|
||||||
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
ft.setMode(modeHang)
|
ft.setMode(modeHang)
|
||||||
|
|
@ -393,7 +485,7 @@ func TestFailedConnectionClosesPortBeforeProbe(t *testing.T) {
|
||||||
// Long interval and a long probe timeout: once the target starts hanging,
|
// Long interval and a long probe timeout: once the target starts hanging,
|
||||||
// a probe cannot complete within the assertion window, so the port closing
|
// a probe cannot complete within the assertion window, so the port closing
|
||||||
// can only be the connection path doing it.
|
// can only be the connection path doing it.
|
||||||
startProxyWith(t, ft, local, 30*time.Second, 5*time.Second)
|
startProxyWith(t, ft, local, 30*time.Second, 5*time.Second, 0)
|
||||||
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
c, err := net.Dial("tcp", local)
|
c, err := net.Dial("tcp", local)
|
||||||
|
|
@ -431,7 +523,7 @@ func TestPortReturnsAfterTransientFailure(t *testing.T) {
|
||||||
local := freePort(t)
|
local := freePort(t)
|
||||||
// Long interval: recovery must come from the probe suspension schedules,
|
// Long interval: recovery must come from the probe suspension schedules,
|
||||||
// not from the next scheduled tick.
|
// not from the next scheduled tick.
|
||||||
startProxyWith(t, ft, local, 30*time.Second, time.Second)
|
startProxyWith(t, ft, local, 30*time.Second, time.Second, 0)
|
||||||
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
waitFor(t, "local port to accept", func() bool { return !localPortRefused(local) })
|
||||||
|
|
||||||
c, err := net.Dial("tcp", local)
|
c, err := net.Dial("tcp", local)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue