Compare commits

...
Author SHA1 Message Date
omarandClaude Opus 5 bcc9df9282 test(wireguard): drive a real AmneziaWG tunnel through ClientBind
release / apk aarch64_cortex-a53 (push) Successful in 3m25s
release / apk x86_64 (push) Successful in 3m14s
release / release apk (push) Successful in 8s
The unit tests pin the reserved-byte gate on each side in isolation, which
would still pass if the two halves disagreed about when to apply it. This wires
two real wireguard-go devices together over loopback UDP through ClientBind on
both ends — the bind the detour path actually uses — configures ranged h1-h4
plus s4 and junk, and asserts an inner IP packet reaches the peer's TUN.

It is red against the unconditional clear and green with the gate, so it covers
the failure the field hit rather than the code we happened to write. Tagged
with_awg, so it runs under the shipped router tag set.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-25 23:25:05 +03:00
omarandClaude Opus 5 ee3641fe45 fix(logsink): collapse interleaved floods, not just consecutive lines
The previous suppression compared each line with the one before it, which the
field never obliges. A dead chain makes the engine cycle the same message
across three outbound tags, so no two identical lines are adjacent: on the
router it produced 854 daemon lines in a ~760-line syslog ring and exactly one
summary, all while claiming "repeated 1 time". The rest of the system's log —
netifd, dnsmasq, the kernel — was evicted anyway.

Track a bounded table of open series keyed by the existing repeat key instead.
The first copy of a key prints; further copies inside its window are counted
whatever arrives in between; the window end emits one summary per key. The
summary now names its message, because several can close at once and "last
message" would simply be false under interleaving.

The table holds 256 keys and evicts the least recently seen, never silently: an
evicted series with a pending count prints its summary on the way out, marked
so the truncation is visible. Close, Reconfigure and any fatal flush every open
series first — a dying daemon may never reach Close.

TestRepeatAlternatingNotSuppressed asserted that A B A B must never be
collapsed. That assertion was the bug. It is replaced by a stronger one: the
messages get separate series, separate summaries and separate counts, so
distinct events still never fold into a single number.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-25 23:23:24 +03:00
omarandClaude Opus 5 439f62238f fix(engine): retire the superseded instance instead of leaving it running
Every config apply built a new box and left the old one alive. The engine's own
log gives it away: inside a single shaterd process, lines carried uptime
counters half an hour apart in the same second, and a live router was found
running four generations at once. A process restart cleared it, so the leak
accrued purely on re-apply.

That is not just wasted memory on a 512 MB box. Each surviving generation keeps
its WireGuard devices up, and two devices sharing one private key evict each
other at the peer — so the leak reproduced the duplicate-device defect between
generations, underneath the deduplication that only reasons about one config.

Retirement now has a hard budget: 5s, which is exactly sing-box's own
C.StopTimeout (past which upstream already calls a stop excessive) and stays
under C.FatalStopTimeout. It is paid after the replacement is serving and only
on an apply that changed something, so a no-op reconcile stays free.

A close that blows the budget is ABANDONED, not waited on, and the apply is
still reported as the success it is — the new box is built, started and
carrying traffic, and failing there would abort the netplane stage and leave a
stale ruleset over a healthy engine. The stuck instance is surfaced through
PendingCloses() into `shaterd status` and the panel, and clears itself if the
shutdown ever completes. Repeated applies over a stuck close no longer stack:
the abandoned generation is remembered, not re-created.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-25 23:23:24 +03:00
omarandClaude Opus 5 d971eb85ee fix(wireguard): stop ClientBind from shredding the AmneziaWG magic header
An AmneziaWG node worked standalone and died the moment it was placed behind
an egress or a chain hop: the handshake completed, the peer answered, and then
not one byte of data ever arrived. The peer never confirmed the session, so it
re-handshook every 15 seconds, forever.

ClientBind cleared bytes 1-3 of every datagram on receive and stamped them on
send, unconditionally. Those bytes are Cloudflare's "reserved" field. They are
also where AmneziaWG puts the upper three bytes of its little-endian uint32
magic header, so zeroing them collapses the value to its low byte, which falls
outside every h1-h4 range and makes the peer classify the packet as an unknown
type and drop it silently.

Handshakes survived because s1/s2 padding pushes their magic past byte 3 — the
clear only scribbled on the random junk prefix. Transport packets have s4 = 0,
so their magic starts at byte 0 and took the hit. That asymmetry is the whole
signature: session up locally, zero data through.

Only the detour path was affected, because Endpoint.Start picks StdNetBind when
the dialer exposes WireGuardControl (no detour) and ClientBind otherwise. The
gate had already landed in StdNetBind; ClientBind was its untouched twin. The
two implement one contract and are now commented as the pair they are, so the
next fix cannot again land on one side only.

Measured on the box: h4 spans 0x60728123-0x60728155, so zeroing bytes 1-3
leaves 35..85 — the captured transport packet began with 56, while a node
without a detour carried a correct 0x6b039798 at the same moment.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-25 23:23:24 +03:00
11 changed files with 2267 additions and 165 deletions
+19 -7
View File
@@ -686,18 +686,30 @@ func (a *Applier) setWarnings(ws []Warning) {
a.stateMu.Unlock()
}
// Warnings returns the normalised warning set from the last successful apply.
// Never nil: an empty slice means "the last apply was clean", which the panel
// must render differently from "no apply has run yet" (Active/Plane cover that).
// Warnings returns the normalised warning set from the last successful apply,
// PLUS whatever is wrong right now that no apply can describe. Never nil: an
// empty slice means "the last apply was clean", which the panel must render
// differently from "no apply has run yet" (Active/Plane cover that).
//
// The live half is currently the engine's abandoned generations
// (engineTeardownWarnings). It is computed at READ time rather than folded into
// lastWarnings on purpose: a superseded box that will not shut down is a
// condition of the process, not a property of a config. Folding it in would make
// it appear only after the NEXT successful apply and then stay published long
// after the shutdown finally completed — reporting a leak that is over, and
// staying silent about one that is not. Read-time means it shows up the instant
// it happens and clears itself the instant it resolves.
func (a *Applier) Warnings() []Warning {
var out []Warning
if a.eng != nil {
out = engineTeardownWarnings(a.eng.PendingCloses())
}
a.stateMu.RLock()
defer a.stateMu.RUnlock()
if a.lastWarnings == nil {
if out == nil && a.lastWarnings == nil {
return []Warning{}
}
out := make([]Warning, len(a.lastWarnings))
copy(out, a.lastWarnings)
return out
return append(out, a.lastWarnings...)
}
// Reconcile re-reads UCI and either tears down (disabled) or re-applies (enabled).
+283
View File
@@ -0,0 +1,283 @@
package apply
// Regression cover for the leaked-engine-generation defect.
//
// Observed on the router: one shaterd process was carrying up to FOUR sing-box
// instances at once. sing-box stamps every log line with the elapsed seconds of
// ITS OWN instance, so the same process printed `ERROR[2015]` and `ERROR[0129]`
// in the same second — two engines half an hour apart in age, both alive, both
// dialling, both holding WireGuard devices built from the same private keys. A
// full daemon stop+start collapsed it back to one generation, which places the
// leak squarely on the config re-apply path rather than on startup.
//
// The tests below pin the two halves of the fix:
//
// TestApplySwapsLeaveExactlyOneEngineGeneration — the healthy path really
// retires the old instance (its listener is provably gone), N times in a row.
// TestStuckEngineCloseDoesNotBlockTheApply — a shutdown that never returns
// is bounded, does not stall the apply, and is REPORTED as a critical warning
// for exactly as long as it is true.
import (
"io"
"net"
"net/netip"
"strconv"
"strings"
"testing"
"time"
C "github.com/sagernet/sing-box/constant"
"github.com/sagernet/sing-box/option"
"github.com/sagernet/sing-box/shater/engine"
"github.com/sagernet/sing/common/json/badoption"
)
// mixedOn builds a minimal but REAL engine config: a mixed inbound bound to
// 127.0.0.1:port plus a direct outbound. Two configs with different ports hash
// differently, so each Apply is a genuine swap rather than a hash-gate no-op —
// and the bound port is the observable that proves whether the old instance
// actually died.
func mixedOn(port uint16) option.Options {
listen := badoption.Addr(netip.MustParseAddr("127.0.0.1"))
return option.Options{
Log: &option.LogOptions{Level: "error"},
Inbounds: []option.Inbound{{
Type: C.TypeMixed,
Tag: "mixed-in",
Options: &option.HTTPMixedInboundOptions{
ListenOptions: option.ListenOptions{Listen: &listen, ListenPort: port},
},
}},
Outbounds: []option.Outbound{{
Type: C.TypeDirect,
Tag: "direct-out",
Options: &option.DirectOutboundOptions{},
}},
}
}
// portFree reports whether 127.0.0.1:port can be bound right now — i.e. whether
// the instance that used to listen there is really gone. Retried briefly because
// a listener is released by Close, not by the return of Close's caller.
func portFree(port uint16) bool {
deadline := time.Now().Add(3 * time.Second)
for {
ln, err := net.Listen("tcp", net.JoinHostPort("127.0.0.1", strconv.Itoa(int(port))))
if err == nil {
_ = ln.Close()
return true
}
if time.Now().After(deadline) {
return false
}
time.Sleep(20 * time.Millisecond)
}
}
// TestApplySwapsLeaveExactlyOneEngineGeneration is the core regression: after N
// sequential applies the process must be carrying ONE engine, not N.
//
// "Carrying" is checked two ways on purpose. Generations() is the engine's own
// accounting (running instance + every retirement still in flight) and would
// catch a retirement that silently never completes. The port check is
// independent of that accounting: if the superseded instance were still alive it
// would still hold its listener, and the bind would fail. A fix that only
// reset a pointer would pass the first check and fail the second.
func TestApplySwapsLeaveExactlyOneEngineGeneration(t *testing.T) {
const (
firstPort = 18801
applies = 5
)
a := New(engine.New(), nil)
t.Cleanup(func() { _ = a.eng.Close() })
for i := 0; i < applies; i++ {
port := uint16(firstPort + i)
changed, err := a.eng.Apply(mixedOn(port))
if err != nil {
t.Fatalf("apply #%d (port %d): %v", i+1, port, err)
}
if !changed {
t.Fatalf("apply #%d: every config here differs, so the swap must be real", i+1)
}
if got := a.eng.Generations(); got != 1 {
t.Fatalf("after apply #%d the process carries %d engine instances, want exactly 1 — "+
"a superseded generation is still alive (this is the four-generations-in-one-process leak)",
i+1, got)
}
if stuck := a.eng.PendingCloses(); len(stuck) != 0 {
t.Fatalf("after apply #%d: %d generation(s) abandoned, want none: %+v", i+1, len(stuck), stuck)
}
if i > 0 {
prev := uint16(firstPort + i - 1)
if !portFree(prev) {
t.Fatalf("after apply #%d the PREVIOUS generation still holds 127.0.0.1:%d — "+
"the old engine was replaced in the field but never actually stopped", i+1, prev)
}
}
// The warning set must stay clean while teardown is healthy: a critical
// warning that cries wolf on every apply is worse than none.
for _, w := range a.Warnings() {
if w.Section == "engine" {
t.Fatalf("after apply #%d a healthy swap produced an engine warning: %+v", i+1, w)
}
}
}
if err := a.eng.Close(); err != nil {
t.Fatalf("close: %v", err)
}
if got := a.eng.Generations(); got != 0 {
t.Fatalf("after Close the process carries %d engine instances, want 0", got)
}
if !portFree(firstPort + applies - 1) {
t.Fatalf("after Close the last generation still holds its listener")
}
}
// hangingCloser returns a box closer that BLOCKS the first close it is handed
// until release() is called, and performs every later close normally. That is the
// shape of the real fault: one subsystem of one generation (a WireGuard endpoint)
// refuses to come down, while the rest of the process is fine.
func hangingCloser() (closer func(io.Closer) error, release func()) {
gate := make(chan struct{})
first := make(chan struct{}, 1)
first <- struct{}{}
return func(c io.Closer) error {
select {
case <-first:
<-gate // the stuck generation: never returns until released
return c.Close() // ...and then really does close, so the port frees
default:
return c.Close()
}
}, func() {
close(gate)
}
}
// TestStuckEngineCloseDoesNotBlockTheApply pins all four requirements of the
// bounded teardown at once:
//
// 1. the apply COMPLETES — a shutdown that never returns must not hold the
// control plane (and therefore the panel) hostage;
// 2. the new engine is running afterwards — fail-closed semantics are unchanged,
// the swap succeeded;
// 3. the abandoned generation is REPORTED as a critical warning, by name, for as
// long as it is still running — it is not silently swallowed;
// 4. generations do not stack: a further apply while the leak persists leaves
// one live instance plus the one abandoned one, not three.
func TestStuckEngineCloseDoesNotBlockTheApply(t *testing.T) {
const budget = 200 * time.Millisecond
restoreBudget := engine.SetCloseBudget(budget)
defer restoreBudget()
a := New(engine.New(), nil)
// Generation 1 comes up with the REAL closer still installed.
if _, err := a.eng.Apply(mixedOn(18811)); err != nil {
t.Fatalf("apply #1: %v", err)
}
closer, release := hangingCloser()
restoreCloser := engine.SetBoxCloser(closer)
released := false
defer func() {
if !released {
release()
}
restoreCloser()
_ = a.eng.Close()
}()
// (1) Generation 2: the retirement of generation 1 will never return.
start := time.Now()
changed, err := a.eng.Apply(mixedOn(18812))
elapsed := time.Since(start)
if err != nil {
t.Fatalf("apply #2 must SUCCEED despite the stuck teardown: %v", err)
}
if !changed {
t.Fatalf("apply #2: expected a real swap")
}
// Generously bounded: the budget plus box.New+Start. The point is that it
// returned at all — before the fix this waited on Box.Close forever.
if elapsed > budget+20*time.Second {
t.Fatalf("apply #2 took %s: a stuck teardown must not stall the apply", elapsed)
}
// (2) fail-closed semantics unchanged: the new engine really is up.
if !a.eng.Running() {
t.Fatalf("apply #2: the new engine must be running")
}
// (3) the leak is visible, named, and critical.
stuck := a.eng.PendingCloses()
if len(stuck) != 1 {
t.Fatalf("PendingCloses() = %+v, want exactly the one abandoned generation", stuck)
}
if stuck[0].Generation != 1 {
t.Errorf("abandoned generation = %d, want 1", stuck[0].Generation)
}
ws := a.Status().Warnings // the exact set `shaterd status` and the panel read
var found *Warning
for i := range ws {
if ws[i].Section == "engine" {
found = &ws[i]
break
}
}
if found == nil {
t.Fatalf("a superseded engine that will not shut down produced NO warning; "+
"Status would show a healthy router: %+v", ws)
}
if found.Severity != SeverityCritical {
t.Errorf("stuck-teardown warning severity = %q, want %q", found.Severity, SeverityCritical)
}
if !strings.Contains(found.Name, "generation 1") {
t.Errorf("stuck-teardown warning must name the generation, got Name=%q", found.Name)
}
if !strings.Contains(found.Message, "STILL RUNNING") {
t.Errorf("stuck-teardown warning must say the instance is still running, got %q", found.Message)
}
if got := a.eng.Generations(); got != 2 {
t.Fatalf("Generations() = %d, want 2 (one live + one abandoned)", got)
}
// (4) another apply while the leak persists must not add a THIRD generation:
// with a generation abandoned the swap goes close-old-then-start-new, so the
// process still holds one live instance plus the one that will not die.
if _, err := a.eng.Apply(mixedOn(18813)); err != nil {
t.Fatalf("apply #3: %v", err)
}
if got := a.eng.Generations(); got != 2 {
t.Fatalf("Generations() = %d after a third apply, want 2 — generations are stacking, "+
"which is exactly the four-live-engines fault", got)
}
if stuck := a.eng.PendingCloses(); len(stuck) != 1 || stuck[0].Generation != 1 {
t.Fatalf("PendingCloses() = %+v, want only the original abandoned generation 1", stuck)
}
// (5) and it CLEARS: when the shutdown finally completes the warning goes away
// on its own. A leak report that outlives the leak trains the operator to
// ignore the panel.
release()
released = true
deadline := time.Now().Add(5 * time.Second)
for len(a.eng.PendingCloses()) > 0 && time.Now().Before(deadline) {
time.Sleep(10 * time.Millisecond)
}
if got := a.eng.PendingCloses(); len(got) != 0 {
t.Fatalf("the finished shutdown is still reported as abandoned: %+v", got)
}
for _, w := range a.Warnings() {
if w.Section == "engine" {
t.Fatalf("the engine warning outlived the leak it describes: %+v", w)
}
}
if got := a.eng.Generations(); got != 1 {
t.Fatalf("Generations() = %d after the stuck shutdown completed, want 1", got)
}
}
+46
View File
@@ -24,7 +24,9 @@ import (
"sort"
"strconv"
"strings"
"time"
"github.com/sagernet/sing-box/shater/engine"
"github.com/sagernet/sing-box/shater/model"
"github.com/sagernet/sing-box/shater/netplane"
)
@@ -269,6 +271,50 @@ func untunnelablePolicyWarnings(g model.Globals, planNotes []string) []Warning {
}
}
// engineTeardownWarnings turns the engine's ABANDONED generations — superseded
// sing-box instances whose shutdown overran the hard close budget and are still
// running inside this process — into operator-facing warnings.
//
// Critical, without hesitation. A leaked generation is not untidiness:
//
// - it still holds its WireGuard devices, and two devices built from the same
// private key evict each other at the peer (one session per public key), so
// the leak reproduces BETWEEN generations exactly the fault
// generate/wgdedup.go removes WITHIN a config — the tunnel flaps and neither
// end can say why;
// - it still holds its outbound connections and keeps probing nodes, so the
// log fills with errors attributed to a config that is no longer applied;
// - on a 512 MiB router each one costs real memory that is never returned.
//
// The generation number is carried in Name so two consecutive status reads can
// tell "the same stuck generation" from "another one just leaked", and the
// elapsed time is in the message because a shutdown at 8s and one at 40 minutes
// are different problems.
func engineTeardownWarnings(stuck []engine.StuckClose) []Warning {
if len(stuck) == 0 {
return nil
}
out := make([]Warning, 0, len(stuck))
for _, s := range stuck {
config := "unknown config"
if len(s.Hash) >= 12 {
config = "config " + s.Hash[:12]
}
out = append(out, Warning{
Severity: SeverityCritical,
Section: "engine",
Name: fmt.Sprintf("generation %d", s.Generation),
Message: fmt.Sprintf(
"a superseded engine instance (%s) has been shutting down for %s and is STILL RUNNING: "+
"it keeps its outbound connections and its WireGuard devices, so it can evict the "+
"live tunnel at the peer and it keeps writing to the log. The current configuration "+
"is applied and running; restart shaterd if this does not clear.",
config, s.Elapsed.Round(time.Second)),
})
}
return out
}
func severityRank(s string) int {
switch s {
case SeverityCritical:
+115 -31
View File
@@ -61,6 +61,25 @@ type Engine struct {
current option.Options // options the running Box was built from
hash string // stable hash of current (hex sha256 of canonical JSON)
// instanceCancel cancels the context THIS box was built on — its own
// cancellable child of e.ctx, one per box (see newBox). box.New does not
// derive a cancellable context of its own, so without this every goroutine
// inside a retired box that waits on ctx.Done() waits forever. Upstream's
// runner does exactly this and calls cancel before Close
// (cmd/sing-box/cmd_run.go). nil when nothing is running.
instanceCancel context.CancelFunc
// instanceGen is the 1-based sequence number of the running box, bumped on
// every adopted swap. It is what names a generation in the log and in
// PendingCloses when a shutdown overruns its budget.
instanceGen uint64
// pending holds the retirements in flight (see teardown.go). Its OWN leaf
// lock, deliberately not mu: PendingCloses is read by the status path, and
// the moment that read matters most is while an apply is holding mu waiting
// out a shutdown that will not finish.
pendingMu sync.Mutex
pending []*pendingClose
// defaultLogWriter, when non-nil, is handed to EVERY box.New this engine
// performs (box.Options.DefaultLogWriter): the daemon points it at its
// long-lived logsink once, and each Apply-swapped box then logs into that
@@ -184,12 +203,18 @@ func (e *Engine) applyLocked(opts option.Options) (bool, error) {
return false, nil
}
// (2) build + validate. box.New constructs and validates every adapter.
nb, err := box.New(box.Options{
Context: e.ctx,
Options: opts,
DefaultLogWriter: e.defaultLogWriter,
})
// (1b) Do not stack generations. A superseded box whose shutdown overran its
// budget is still holding its outbound connections and its WireGuard devices;
// building another one on top of it is how one process ends up carrying four
// live generations. Give any abandoned teardown one more budget to finish
// BEFORE we create anything (see teardown.go awaitAbandonedLocked). Placed
// after the hash gate on purpose: a no-op reconcile — cron, every minute —
// must stay free.
stuck := e.awaitAbandonedLocked()
// (2) build + validate on its OWN cancellable context. box.New constructs and
// validates every adapter.
nb, nbCancel, err := e.newBox(opts)
if err != nil {
// Validation failed: keep the running instance, do not swap.
return false, E.Cause(err, "create instance")
@@ -211,8 +236,15 @@ func (e *Engine) applyLocked(opts option.Options) (bool, error) {
// close-old-then-start-new DIRECTLY — skipping the stall. (box.New above
// already validated opts, so we never tear the old box down for a config
// that would fail to build.)
if e.instance != nil && sharesCacheFileLock(e.current, opts) {
return e.closeOldThenStart(nb, opts, newHash)
//
// A generation that is STILL abandoned after the wait above forces the same
// path for a different reason: start-new-first would put a second LIVE box
// alongside a third that refuses to die, all three contending for the same
// tproxy port, the same cache file and — the expensive one — the same
// WireGuard private keys. Closing the current box first keeps the process to
// at most one live instance plus the abandoned one.
if e.instance != nil && (stuck > 0 || sharesCacheFileLock(e.current, opts)) {
return e.closeOldThenStart(nb, nbCancel, opts, newHash)
}
// (3b) start-new-first (zero-downtime when there is no resource conflict).
@@ -223,25 +255,59 @@ func (e *Engine) applyLocked(opts option.Options) (bool, error) {
// close-old-then-start-new then. Any OTHER Start failure keeps today's
// behavior: discard the new box, keep the old one running.
if e.instance == nil || !isSwapConflict(err) {
nbCancel()
_ = nb.Close()
return false, E.Cause(err, "start instance")
}
return e.closeOldThenStart(nb, opts, newHash)
return e.closeOldThenStart(nb, nbCancel, opts, newHash)
}
// Swap succeeded. The previously running config becomes last-good.
old := e.instance
if old != nil {
if e.instance != nil {
e.lastGood = e.current
e.hasLastGood = true
_ = old.Close()
// Retire the old generation under the hard budget. Its error — including
// ErrCloseTimeout — is deliberately NOT returned: the new box is started
// and carrying traffic, so this apply SUCCEEDED, and failing it here would
// abort the caller's netplane stage and leave a stale ruleset loaded over a
// perfectly healthy engine. An abandoned generation is surfaced through
// PendingCloses() (critical warning in `shaterd status` and the panel) and
// an ERROR line naming it — visible, but not mistaken for a failed apply.
_ = e.retireLocked(e.instance, e.instanceCancel, e.instanceGen, e.hash)
}
e.instance = nb
e.current = opts
e.hash = newHash
e.adoptLocked(nb, nbCancel, opts, newHash)
return true, nil
}
// newBox builds a box on its OWN cancellable child of the engine context and
// returns the cancel alongside it. Every goroutine the box starts inherits that
// context, so cancelling it is what unwinds the ones Close does not reach; see
// teardown.go for why the shared, never-cancelled context was the defect.
func (e *Engine) newBox(opts option.Options) (*box.Box, context.CancelFunc, error) {
ctx, cancel := context.WithCancel(e.ctx)
b, err := box.New(box.Options{
Context: ctx,
Options: opts,
DefaultLogWriter: e.defaultLogWriter,
})
if err != nil {
cancel()
return nil, nil, err
}
return b, cancel, nil
}
// adoptLocked installs a started box as THE running instance and gives it the
// next generation number. Caller holds e.mu and has already retired whatever was
// running before.
func (e *Engine) adoptLocked(b *box.Box, cancel context.CancelFunc, opts option.Options, hash string) {
e.instanceGen++
e.instance = b
e.instanceCancel = cancel
e.current = opts
e.hash = hash
}
// closeOldThenStart is the close-old-then-start-new swap. It is taken both
// proactively (the incoming config shares the running box's cache_file lock, so
// start-new-first cannot work) and reactively (start-new-first hit a swap
@@ -256,27 +322,36 @@ func (e *Engine) applyLocked(opts option.Options) (bool, error) {
// already validated by the caller's box.New, so the rebuild below is expected to
// succeed; the restore path guards the rare case it does not. The caller holds
// e.mu.
func (e *Engine) closeOldThenStart(discard *box.Box, opts option.Options, newHash string) (bool, error) {
func (e *Engine) closeOldThenStart(discard *box.Box, discardCancel context.CancelFunc, opts option.Options, newHash string) (bool, error) {
// (a) Drop the pre-built box (a Start-failed box cannot be restarted, and the
// proactively-built one must not hold the cache_file lock while we rebuild).
// It was never adopted, so it is not a generation — close it inline, but
// cancel its context first exactly like a retired one.
if discardCancel != nil {
discardCancel()
}
if discard != nil {
_ = discard.Close()
}
// (b) Free the port + cache_file lock by closing the old instance. Remember
// its config so we can restore it if the fresh box cannot come up.
// (b) Free the port + cache_file lock by closing the old instance, under the
// hard budget. Remember its config so we can restore it if the fresh box
// cannot come up. An overrun here is reported by retireLocked and tracked in
// PendingCloses; we still proceed, because the resources it was supposed to
// free are exactly what the operator is waiting on.
prevOpts := e.current
prevHash := e.hash
prevLastGood := e.lastGood
prevHasLastGood := e.hasLastGood
_ = e.instance.Close()
e.instance = nil
_ = e.retireLocked(e.instance, e.instanceCancel, e.instanceGen, prevHash)
e.instance, e.instanceCancel = nil, nil
// (c) Build a FRESH box for opts (the discarded one cannot be reused).
nb2, err := box.New(box.Options{Context: e.ctx, Options: opts, DefaultLogWriter: e.defaultLogWriter})
nb2, cancel2, err := e.newBox(opts)
if err == nil {
err = nb2.Start()
if err != nil {
cancel2()
_ = nb2.Close()
}
}
@@ -284,28 +359,25 @@ func (e *Engine) closeOldThenStart(discard *box.Box, opts option.Options, newHas
// (d) Success: the old config we just closed becomes last-good.
e.lastGood = prevOpts
e.hasLastGood = true
e.instance = nb2
e.current = opts
e.hash = newHash
e.adoptLocked(nb2, cancel2, opts, newHash)
return true, nil
}
// (e) The fresh box could not come up and the old one is already closed —
// interception is currently down. Try to RESTORE the previous config so we
// do not leave the tunnel dead.
rb, rerr := box.New(box.Options{Context: e.ctx, Options: prevOpts, DefaultLogWriter: e.defaultLogWriter})
rb, rcancel, rerr := e.newBox(prevOpts)
if rerr == nil {
rerr = rb.Start()
if rerr != nil {
rcancel()
_ = rb.Close()
}
}
if rerr == nil {
// Old config restored: keep current/hash/last-good exactly as they were
// (do NOT advance them). Report that opts was not applied.
e.instance = rb
e.current = prevOpts
e.hash = prevHash
e.adoptLocked(rb, rcancel, prevOpts, prevHash)
e.lastGood = prevLastGood
e.hasLastGood = prevHasLastGood
return false, E.Cause(err, "start instance (config not applied; previous config restored)")
@@ -314,7 +386,7 @@ func (e *Engine) closeOldThenStart(discard *box.Box, opts option.Options, newHas
// Restore ALSO failed: the engine is now STOPPED. e.instance stays nil (never
// pointing at a closed box). The fail-closed nft kill-switch keeps the LAN
// safe (no unproxied leak) even though interception is down.
e.instance = nil
e.instance, e.instanceCancel = nil, nil
e.hash = ""
return false, E.Cause(E.Errors(err, rerr), "start instance failed and could not restore previous config; engine stopped")
}
@@ -410,14 +482,26 @@ func (e *Engine) HasLastGood() bool {
}
// Close stops the running instance, if any. It is idempotent.
//
// Unlike the swap path, this one DOES return ErrCloseTimeout: here the shutdown
// is the whole operation, so "it did not stop" is the result, not a footnote. The
// caller (apply.Teardown) records it as the teardown's error while still
// completing the netplane teardown — the data plane must come down even when a
// box will not.
func (e *Engine) Close() error {
e.mu.Lock()
defer e.mu.Unlock()
if e.instance == nil {
// Nothing running, but a previously abandoned generation may still be:
// give it a last budget so a teardown followed by a restart does not carry
// the leak across.
if e.awaitAbandonedLocked() > 0 {
return ErrCloseTimeout
}
return nil
}
err := e.instance.Close()
e.instance = nil
err := e.retireLocked(e.instance, e.instanceCancel, e.instanceGen, e.hash)
e.instance, e.instanceCancel = nil, nil
e.hash = ""
return err
}
+344
View File
@@ -0,0 +1,344 @@
package engine
// Bounded, deterministic teardown of a SUPERSEDED sing-box instance.
//
// # The defect this file exists for
//
// sing-box's Box has no live reload, so every applied config change builds a
// fresh Box and retires the old one (see engine.go). Retiring it used to be one
// line — `_ = old.Close()` — and that line had three separate problems, all of
// which had to be true at once for the observed failure:
//
// 1. NO CANCELLATION. Upstream's own runner gives every Box its OWN cancellable
// context and calls cancel() BEFORE Close (cmd/sing-box/cmd_run.go:136 and
// :188-190). This engine handed EVERY box.New the one shared, never-cancelled
// context built in New (engine.go), and box.New does not derive a cancellable
// child of what it is given (box.go: `ctx := options.Context` and nothing
// else). So every goroutine inside a retired box that would have stopped on
// context cancellation simply never stopped, and only what each adapter's
// Close() explicitly tears down actually went away.
//
// 2. NO TIME BUDGET. Box.Close walks its subsystems SEQUENTIALLY and waits for
// each one forever; the only thing watching the clock is a taskmonitor that
// prints "close endpoint/wireguard[...] take too much time to finish!" after
// C.StopTimeout and then keeps waiting anyway (common/taskmonitor/monitor.go).
// A wireguard endpoint that will not come down therefore stalls the rest of
// the close list, and every subsystem AFTER it in the walk is never reached.
//
// 3. NO EVIDENCE. The error was discarded (`_ =`) and the pointer to the old box
// was overwritten in the same breath, so after the swap nothing in the process
// could tell — or even ask — whether the previous generation had actually
// died. On the router this accumulated: up to four generations logging side by
// side inside one shaterd process, each with its own outbound connections and,
// worse, its own WireGuard devices. Two devices sharing a private key evict
// each other at the peer, which is exactly the fault generate/wgdedup.go
// removes WITHIN a config — reintroduced here BETWEEN generations.
//
// # The contract now
//
// - Every box gets its own cancellable context, cancelled before Close.
// - Close runs against a HARD budget (closeBudget). When it expires the apply
// moves on: the new box is already started and serving, and blocking the
// control plane on a shutdown that is not going to finish would only add an
// unresponsive panel to the problem.
// - An overrun is LOUD, not swallowed: an ERROR line naming the generation, and
// a critical entry in the operator-facing warning set for as long as the
// abandoned generation is still running (see PendingCloses).
// - Generations do not stack: before an apply builds another instance it gives
// any abandoned teardown one more budget to finish, and if one is still alive
// the swap is forced onto the close-old-then-start-new path so the process
// never holds two LIVE boxes on top of an abandoned one.
import (
"errors"
"io"
"sync"
"time"
box "github.com/sagernet/sing-box"
)
// ErrCloseTimeout reports that a superseded instance did not finish shutting
// down inside the close budget and has been ABANDONED (it may still be running).
//
// It is deliberately NOT propagated out of Apply. The swap it accompanies
// succeeded — the new box is built, started and carrying traffic — and returning
// a failure there would make the caller abort the netplane stage and, with no
// engine fault to point at, leave a stale ruleset loaded over a healthy engine.
// The fact is surfaced through PendingCloses() instead, which reaches
// `shaterd status` and the panel and clears itself the moment the shutdown
// finally completes. Close()/Teardown DO return it: there the failure to stop is
// the whole point of the operation.
var ErrCloseTimeout = errors.New("engine instance did not shut down within the close budget")
// defaultCloseBudget is the HARD wall-clock budget for retiring one superseded
// box.
//
// Five seconds, for three reasons that all point at the same number:
//
// - it is exactly sing-box's own C.StopTimeout — the threshold at which the
// engine itself declares a single lifecycle stop to be taking too long. A
// close that blows past the budget upstream considers excessive is by
// definition not a slow close, it is a stuck one;
// - it stays under C.FatalStopTimeout (10s), which is when upstream's CLI gives
// up and calls the process unclosable. We want to have already reacted by
// then;
// - it is paid AFTER the replacement box is started and serving, and only on an
// apply that really changed something (the hash gate makes a no-op reconcile
// free), so the worst case is five seconds added to one apply — not to the
// every-minute cron reconcile, and never to a status poll.
const defaultCloseBudget = 5 * time.Second
// closeBudget/closeBoxFn are package-wide seams. Production never touches them;
// tests in this package and in shater/apply install a short budget and a
// deliberately hung closer to exercise the abandoned-generation path without
// standing up a box that really refuses to die.
var (
seamMu sync.RWMutex
closeBudget = defaultCloseBudget
closeBoxFn = func(c io.Closer) error { return c.Close() }
)
// SetCloseBudget overrides the close budget and returns a function restoring the
// previous value. TEST SEAM — production uses defaultCloseBudget.
func SetCloseBudget(d time.Duration) (restore func()) {
seamMu.Lock()
prev := closeBudget
closeBudget = d
seamMu.Unlock()
return func() {
seamMu.Lock()
closeBudget = prev
seamMu.Unlock()
}
}
// SetBoxCloser overrides HOW a superseded instance is closed and returns a
// function restoring the previous closer. TEST SEAM — production calls
// Box.Close. A closer that never returns is how the abandoned-generation path is
// tested.
func SetBoxCloser(fn func(io.Closer) error) (restore func()) {
seamMu.Lock()
prev := closeBoxFn
closeBoxFn = fn
seamMu.Unlock()
return func() {
seamMu.Lock()
closeBoxFn = prev
seamMu.Unlock()
}
}
func currentCloseBudget() time.Duration {
seamMu.RLock()
defer seamMu.RUnlock()
return closeBudget
}
func currentBoxCloser() func(io.Closer) error {
seamMu.RLock()
defer seamMu.RUnlock()
return closeBoxFn
}
// pendingClose tracks ONE retirement in flight. It lives in Engine.pending from
// the moment the retirement starts until Close returns, and `abandoned` records
// whether the budget expired while it was still running — i.e. whether this
// generation is a leak the operator must be told about, or merely a close that
// is in progress under a lock the caller is holding anyway.
type pendingClose struct {
gen uint64
hash string // config hash the retired box was built from
since time.Time
done chan struct{} // closed when the underlying Close finally returns
abandoned bool // guarded by Engine.pendingMu
finished bool // guarded by Engine.pendingMu
}
// StuckClose is the read-side view of a superseded generation that overran the
// close budget and is STILL running inside this process.
type StuckClose struct {
// Generation is the 1-based sequence number of the box that will not die.
// It matches nothing in the sing-box log by itself, but it lets two status
// reads tell "the same old generation" from "another one just leaked".
Generation uint64
// Hash is the config hash the leaked box was built from — the same value
// `shaterd status` reported as `hash` while that config was the running one.
// It is what connects "an engine is stuck" to WHICH configuration is stuck.
Hash string
// Since is when its shutdown was started.
Since time.Time
// Elapsed is how long it has been shutting down, as of the read.
Elapsed time.Duration
}
// retireLocked tears down a superseded box under the close budget. The caller
// holds e.mu and must clear e.instance itself.
//
// Order matters and mirrors upstream: cancel the box's context FIRST so every
// goroutine keyed on it unwinds, and only then call Close, which is what actually
// releases the listeners, the WireGuard devices and the cache-file lock.
//
// Returns nil on a clean shutdown, ErrCloseTimeout when the budget expired (the
// close keeps running in the background and is reaped when it finishes), or the
// close's own error.
func (e *Engine) retireLocked(b *box.Box, cancel func(), gen uint64, hash string) error {
if cancel != nil {
cancel()
}
if b == nil {
return nil
}
pc := &pendingClose{gen: gen, hash: hash, since: time.Now(), done: make(chan struct{})}
e.pendingMu.Lock()
e.pending = append(e.pending, pc)
e.pendingMu.Unlock()
closeFn := currentBoxCloser()
errc := make(chan error, 1)
go func() {
err := closeFn(b)
e.pendingMu.Lock()
pc.finished = true
wasAbandoned := pc.abandoned
e.removePendingLocked(pc)
e.pendingMu.Unlock()
if wasAbandoned && e.log != nil {
// The counterpart of the ERROR below: an operator who saw the leak
// reported must be able to see it resolve without restarting anything.
e.log.Warn("engine generation ", gen, " finally finished shutting down after ",
time.Since(pc.since).Round(time.Millisecond), "; it is no longer running")
}
close(pc.done)
errc <- err
}()
budget := currentCloseBudget()
timer := time.NewTimer(budget)
defer timer.Stop()
select {
case err := <-errc:
return err
case <-timer.C:
}
// The budget expired. Claim the generation as abandoned — unless the close
// happened to land in the same instant, in which case there is nothing to
// report and we take its real result.
e.pendingMu.Lock()
abandoned := !pc.finished
if abandoned {
pc.abandoned = true
}
e.pendingMu.Unlock()
if !abandoned {
return <-errc
}
if e.log != nil {
e.log.Error("engine generation ", gen, " (config ", shortHash(hash), ") did NOT shut down within ", budget,
" and has been ABANDONED: it may still hold its outbound connections and its ",
"WireGuard devices (two devices with the same private key evict each other at ",
"the peer), and it keeps writing to this log. The new configuration is running; ",
"restart shaterd if this generation never clears.")
}
return ErrCloseTimeout
}
// shortHash renders a config hash the way the log and the panel both need it:
// enough to identify the configuration, short enough to read in a syslog line.
func shortHash(h string) string {
if len(h) < 12 {
return "unknown"
}
return h[:12]
}
// removePendingLocked drops pc from the pending list. Caller holds pendingMu.
func (e *Engine) removePendingLocked(pc *pendingClose) {
for i, p := range e.pending {
if p == pc {
e.pending = append(e.pending[:i], e.pending[i+1:]...)
return
}
}
}
// pendingSnapshot copies the in-flight retirements. Leaf lock only.
func (e *Engine) pendingSnapshot() []*pendingClose {
e.pendingMu.Lock()
defer e.pendingMu.Unlock()
return append([]*pendingClose(nil), e.pending...)
}
// PendingCloses lists the superseded generations that overran their close budget
// and are still running. Empty is the healthy answer.
//
// It takes ONLY the pending leaf lock — never e.mu — so the panel and
// `shaterd status` can report a leaking teardown even while the apply that
// produced it is still holding the engine mutex. That is deliberate: the one
// moment this information matters most is while an apply is stalled behind a
// shutdown that will not finish.
func (e *Engine) PendingCloses() []StuckClose {
now := time.Now()
e.pendingMu.Lock()
defer e.pendingMu.Unlock()
out := make([]StuckClose, 0, len(e.pending))
for _, p := range e.pending {
if !p.abandoned {
continue
}
out = append(out, StuckClose{
Generation: p.gen,
Hash: p.hash,
Since: p.since,
Elapsed: now.Sub(p.since),
})
}
return out
}
// Generations reports how many sing-box instances this process is carrying: the
// running one (0 or 1) plus every superseded generation whose shutdown has not
// finished. ONE is the healthy answer for a started engine; anything above it is
// the leak this file exists to make impossible to hide.
func (e *Engine) Generations() int {
e.mu.Lock()
live := 0
if e.instance != nil {
live = 1
}
e.mu.Unlock()
e.pendingMu.Lock()
defer e.pendingMu.Unlock()
return live + len(e.pending)
}
// awaitAbandonedLocked gives every already-abandoned generation ONE more close
// budget to finish, and returns how many are still running afterwards. The caller
// holds e.mu and is about to build another instance.
//
// This is what keeps generations from stacking. Without it, a box that will not
// come down means the NEXT apply quietly adds a third instance to the process,
// and the one after that a fourth — which is precisely the accumulation observed
// on the router. The wait is bounded by one budget for the whole set (not per
// entry), so a permanently stuck generation costs a single extra budget on an
// apply that actually changes the config, and nothing at all on the every-minute
// no-op reconcile, which never gets past the hash gate.
func (e *Engine) awaitAbandonedLocked() int {
pending := e.pendingSnapshot()
if len(pending) == 0 {
return 0
}
deadline := time.NewTimer(currentCloseBudget())
defer deadline.Stop()
for _, p := range pending {
select {
case <-p.done:
case <-deadline.C:
return len(e.PendingCloses())
}
}
return len(e.PendingCloses())
}
+247
View File
@@ -0,0 +1,247 @@
// lx: pins the invariant that a chain copy of an AmneziaWG node carries the
// SAME obfuscation parameters as the base node, field for field.
//
// Context: on the router a node placed behind an egress hop
// (egress:ewan -> node:awgout, tag chain-<name>-h1) handshook forever and never
// passed a transport packet, while the same node standalone was healthy. One
// candidate explanation was that rebuildNode loses AWG params on the copy — it
// does not (both the base and the copy go through the same
// builder.wireguardEndpoint / amneziaOptions, the copy differing only in Tag and
// DialerOptions.Detour), and this test holds that line. The real cause was the
// bind swap the detour triggers: a detour makes Endpoint.Start pick ClientBind
// instead of conn.StdNetBind, and ClientBind unconditionally overwrote bytes 1-3
// of every datagram, shredding the h4 magic header of transport packets. See
// transport/wireguard/client_bind.go (hasReserved) and its regression tests.
//
// Honest note: this test passes both before and after that fix — it is a pin on
// a path that was never broken, not the reproducer for the bug.
//
// Unlike generate_test.go this file is NOT Linux-gated: it stops at
// GenerateWithWarnings and never calls engine.Apply / box.New, so it needs no
// routing_mark validation and runs on every platform.
package generate
import (
"encoding/base64"
"fmt"
"net/url"
"reflect"
"testing"
"github.com/sagernet/sing-box/option"
"github.com/sagernet/sing-box/shater/model"
)
// awgKey returns a valid 32-byte base64 WireGuard key seeded by fill. Local to
// this file so it does not depend on the Linux-only suite's helpers.
func awgKey(fill byte) string {
b := make([]byte, 32)
for i := range b {
b[i] = fill + byte(i)
}
return base64.StdEncoding.EncodeToString(b)
}
// TestChainCopyPreservesAmneziaWGOptions builds a chain whose terminal hop is an
// AmneziaWG node and asserts the chain copy's AmneziaWGOptions equals the base
// endpoint's, comparing every field of the struct (so a field added later
// without being mapped in amneziaOptions is caught here too).
func TestChainCopyPreservesAmneziaWGOptions(t *testing.T) {
priv := awgKey(1)
pub := awgKey(9)
// Every knob the option struct carries that the share-link can express:
// jc/jmin/jmax, s1-s3, ranged h1-h4 (AWG 2.0), i1-i5. s4 is deliberately left
// unset (0) — that is the real-world shape in which the transport magic lands
// in bytes 0-3 and the ClientBind bug bit.
uri := fmt.Sprintf(
"awg://%s@203.0.113.10:51820?publickey=%s&address=10.13.13.2/32&allowedips=0.0.0.0/0"+
"&jc=4&jmin=40&jmax=70&s1=86&s2=57&s3=13"+
"&h1=1618116899-1618116949&h2=1795397486-1795397536&h3=3333333333&h4=1618116899-1618116949"+
"&i1=%s&i2=%s&i3=%s&i4=%s&i5=%s#awgout",
url.QueryEscape(priv), url.QueryEscape(pub),
url.QueryEscape("<b 0xf0>"), url.QueryEscape("<c>"), url.QueryEscape("<t>"),
url.QueryEscape("<r 10>"), url.QueryEscape("<b 0xab>"),
)
node := model.Node{Name: "awgout", Enabled: true, URI: uri}
// Two SEPARATE configs, not one. wgdedup refuses to materialise a WireGuard
// node twice from one private key (two devices sharing a key evict each other's
// session), so a single config can hold either the base endpoint or the chain
// copy — never both. That is also why on the router this node exists only as
// "chain-<name>-h1". The invariant under test is therefore cross-config: the
// same node, referenced directly vs referenced through a chain, must yield the
// same AmneziaWG parameters.
direct := &model.Model{
Globals: model.DefaultGlobals(),
Nodes: []model.Node{node},
Rules: []model.Rule{
{Name: "default", Enabled: true, Order: 100, Target: "node:awgout"},
},
}
chained := &model.Model{
Globals: model.DefaultGlobals(),
Nodes: []model.Node{node},
Egresses: []model.Egress{
{Name: "ewan", Type: "interface", Interface: "wan"},
},
Chains: []model.Chain{
{Name: "viaewan", Hops: []string{"egress:ewan", "node:awgout"}},
},
Rules: []model.Rule{
{Name: "default", Enabled: true, Order: 100, Target: "chain:viaewan"},
},
}
directOpts, warns, err := GenerateWithWarnings(direct)
if err != nil {
t.Fatalf("GenerateWithWarnings(direct): %v (warnings: %v)", err, warns)
}
chainedOpts, chainWarns, err := GenerateWithWarnings(chained)
if err != nil {
t.Fatalf("GenerateWithWarnings(chained): %v (warnings: %v)", err, chainWarns)
}
base, ok := awgOptionsForTag(directOpts, "awgout")
if !ok {
t.Fatalf("base endpoint %q not found; tags = %v (warnings: %v)",
"awgout", endpointTags(directOpts), warns)
}
if !base.IsSet() {
t.Fatalf("base endpoint has no AmneziaWG params: %+v", base)
}
// Absolute expectation, not just parity. A "copy == base" assertion alone is
// satisfied when BOTH lose a field (they share amneziaOptions), so pin the
// literal values the share-link carries. Adding a field to
// option.AmneziaWGOptions without mapping it in amneziaOptions fails the
// exhaustiveness check below.
want := option.AmneziaWGOptions{
Jc: 4, Jmin: 40, Jmax: 70,
S1: 86, S2: 57, S3: 13, S4: 0,
H1: "1618116899-1618116949",
H2: "1795397486-1795397536",
H3: "3333333333",
H4: "1618116899-1618116949",
I1: "<b 0xf0>", I2: "<c>", I3: "<t>", I4: "<r 10>", I5: "<b 0xab>",
}
if base != want {
t.Fatalf("base AmneziaWG params lost/garbled in parsing:\n got %+v\nwant %+v", base, want)
}
// Exhaustiveness: every field of the struct must be exercised above, so a
// newly added knob cannot slip through unmapped and unnoticed.
assertAWGFieldsCovered(t, want)
// The chain hop copy: buildHopWrapper tags it chain-<chain>-h<idx>. Find it as
// "the wireguard endpoint in the chained config" rather than hardcoding the
// index, so a change in hop numbering does not silently no-op this test.
var (
copyTag string
copyOptions option.AmneziaWGOptions
found bool
)
for _, endpoint := range chainedOpts.Endpoints {
wg, isWG := endpoint.Options.(*option.WireGuardEndpointOptions)
if !isWG {
continue
}
copyTag, copyOptions, found = endpoint.Tag, wg.AmneziaWGOptions, true
break
}
if !found {
t.Fatalf("no chain copy endpoint emitted; endpoint tags = %v (warnings: %v)",
endpointTags(chainedOpts), chainWarns)
}
if copyTag == "awgout" {
t.Fatalf("expected a per-hop chain copy tag, got the base tag %q", copyTag)
}
// Field-by-field, via reflection: any field of AmneziaWGOptions the copy fails
// to carry is named explicitly rather than hidden behind one struct diff.
baseValue := reflect.ValueOf(base)
copyValue := reflect.ValueOf(copyOptions)
for i := 0; i < baseValue.NumField(); i++ {
field := baseValue.Type().Field(i)
want := baseValue.Field(i).Interface()
got := copyValue.Field(i).Interface()
if !reflect.DeepEqual(want, got) {
t.Errorf("chain copy %q lost AmneziaWG field %s: got %#v, want %#v",
copyTag, field.Name, got, want)
}
}
if !t.Failed() && base != copyOptions {
t.Fatalf("chain copy %q AmneziaWGOptions differ from base: %+v vs %+v",
copyTag, copyOptions, base)
}
// The copy must additionally differ from the base in exactly the way the chain
// intends: it detours through the egress hop, the base does not.
baseDialer, okBase := dialerForTag(directOpts, "awgout")
copyDialer, okCopy := dialerForTag(chainedOpts, copyTag)
if !okBase || !okCopy {
t.Fatalf("dialer options missing: base=%v copy=%v", okBase, okCopy)
}
if baseDialer.Detour != "" {
t.Errorf("base endpoint must dial directly, got Detour=%q", baseDialer.Detour)
}
if copyDialer.Detour == "" {
t.Errorf("chain copy %q must carry the egress hop as Detour, got empty", copyTag)
}
}
// assertAWGFieldsCovered fails if any field of option.AmneziaWGOptions is left
// at its zero value in the expectation, other than the ones deliberately unset:
//
// S4 — kept 0 on purpose; that is the shape in which the transport magic
// lands in bytes 0-3, i.e. the configuration that broke on the router.
// Id/Ip/Ib — WireSock masquerade sugar, mutually exclusive with an explicit I1
// (device_awg.go rejects the combination), and I1 is set here.
//
// The point is that adding a knob to the struct without extending this test
// turns into a failure here rather than silent non-coverage.
func assertAWGFieldsCovered(t *testing.T, want option.AmneziaWGOptions) {
t.Helper()
deliberatelyUnset := map[string]bool{"S4": true, "Id": true, "Ip": true, "Ib": true}
value := reflect.ValueOf(want)
for i := 0; i < value.NumField(); i++ {
name := value.Type().Field(i).Name
if deliberatelyUnset[name] {
continue
}
if value.Field(i).IsZero() {
t.Errorf("AmneziaWGOptions.%s is not exercised by this test "+
"(add it to the share-link and to `want`, or to deliberatelyUnset)", name)
}
}
}
// awgOptionsForTag returns the AmneziaWGOptions of the wireguard endpoint tagged
// tag.
func awgOptionsForTag(opts option.Options, tag string) (option.AmneziaWGOptions, bool) {
for _, endpoint := range opts.Endpoints {
if endpoint.Tag != tag {
continue
}
wg, ok := endpoint.Options.(*option.WireGuardEndpointOptions)
if !ok {
return option.AmneziaWGOptions{}, false
}
return wg.AmneziaWGOptions, true
}
return option.AmneziaWGOptions{}, false
}
// dialerForTag returns the DialerOptions of the wireguard endpoint tagged tag.
func dialerForTag(opts option.Options, tag string) (option.DialerOptions, bool) {
for _, endpoint := range opts.Endpoints {
if endpoint.Tag != tag {
continue
}
wg, ok := endpoint.Options.(*option.WireGuardEndpointOptions)
if !ok {
return option.DialerOptions{}, false
}
return wg.DialerOptions, true
}
return option.DialerOptions{}, false
}
+260 -92
View File
@@ -16,11 +16,15 @@
// answer "give me the last day" — the prefix is what makes the download
// ranges real. UTC only: the binary ships without tzdata, so any local-zone
// rendering would be a fiction.
// - collapses RUNS of the same message into "last message repeated N times"
// (see repeatKey / repeatWindow): a broken outbound makes the engine repeat
// one line hundreds of times a minute, which evicts the whole rest of the
// router's syslog ring buffer within minutes. Suppression applies to both
// halves identically; fatal/panic lines are never suppressed.
// - collapses REPEATS of the same message into "repeated N times: <message>"
// (see repeatKey / repeatWindow / maxRepeatKeys): a broken outbound makes
// the engine repeat one line hundreds of times a minute, which evicts the
// whole rest of the router's syslog ring buffer within minutes. The repeats
// do NOT have to be adjacent — a flood usually interleaves a handful of
// messages (one per broken chain), so the sink keeps a small table of the
// messages seen recently instead of comparing with the previous line only.
// Suppression applies to both halves identically; fatal/panic lines are
// never suppressed.
// - fans it out according to Config: to the REAL os.Stderr (procd relays fd2
// to syslog/logread) when ToSyslog, and to a size-capped, 2-segment rotated
// file when ToFile. Both off => the line is dropped — that IS the "fully
@@ -47,6 +51,7 @@ package logsink
import (
"bytes"
"container/list"
"fmt"
"io"
"os"
@@ -80,21 +85,44 @@ const (
// that grows to its full cap leaves the rootfs this much headroom.
diskFloorBytes = 4 << 20
// repeatWindow bounds how long a run of identical messages may stay silent:
// once suppression starts, a "last message repeated N times" summary is
// emitted every window for as long as the run continues, and the counter
// restarts.
// repeatWindow bounds how long copies of one message may stay silent: the
// first occurrence is printed and opens a window; every further copy of THAT
// message inside the window is swallowed and reported by a
// "repeated N times: <message>" summary when the window ends. A message that
// is still flooding keeps its window (one summary per window, counter
// restarting); a message that went quiet is forgotten, so it prints in full
// the next time it happens.
//
// Why 5s: the observed flood (a dead chain -> "WireGuard is not ready yet")
// runs at ~1 line/s, i.e. ~370 lines per 6 minutes, while the router's
// syslog ring holds ~760 lines total — one faulty outbound erases every
// other subsystem's history, including our own startup lines. A 5s window
// turns that into ~12 lines/min (~30x less) — small enough that the ring
// survives a long outage, short enough that an operator tailing `logread -f`
// sees the fault acknowledged within 5 seconds and keeps seeing it, so a
// standing problem never looks like a frozen log. A longer window (30-60s)
// would hide the ongoing-ness; a shorter one would not relieve the ring.
// Why 5s: the observed flood (dead chains -> "WireGuard is not ready yet")
// runs at ~7 lines/s across three chain tags — 436 lines in one minute —
// while the router's syslog ring holds ~760 lines total: one faulty
// subscription erases every other subsystem's history, including our own
// startup lines, in under two minutes. A 5s window turns that into one
// summary per message per window (~36 lines/min instead of 436) — short
// enough that an operator tailing `logread -f` sees the fault acknowledged
// within 5 seconds and keeps seeing it, so a standing problem never looks
// like a frozen log. A longer window (30-60s) would hide the ongoing-ness; a
// shorter one would not relieve the ring.
repeatWindow = 5 * time.Second
// maxRepeatKeys bounds the table of recently seen messages. Suppression must
// work across INTERLEAVED messages, so the sink cannot keep just the last
// key — but the keys come from the engine and are as varied as its log, so
// the table must not be allowed to grow with them.
//
// Why 256: the floods worth collapsing are per-outbound or per-chain, and a
// pathological config on this box has a few hundred nodes of which only the
// broken handful actually log; 256 distinct messages in flight covers that
// with room to spare, while a table this size is trivial next to the
// process — entries hold the log line itself, ~150 B typically (~40 KiB
// total) and 8 KiB at the absolute worst (maxPartialLine), i.e. ~2 MiB even
// in the case that cannot really happen.
//
// Eviction is least-recently-seen: the entry that has gone longest without a
// copy is the one least likely to be flooding. Eviction is never silent —
// an evicted entry that had swallowed copies prints its summary on the way
// out (marked "repeat table full"), so a counter is never simply dropped.
maxRepeatKeys = 256
)
// FilePath returns the log-file location for the given persistence choice.
@@ -174,11 +202,13 @@ type Sink struct {
suspended bool // persistent-path disk guard tripped
lastProbe time.Time // last disk-free probe
// repeat suppression (emitLocked / flushRepeatsLocked)
lastKey string // repeatKey of the last PRINTED line
hasLast bool // a lastKey exists (distinguishes "" from "unset")
repeats int // identical lines swallowed since the last summary
repeatTimer *time.Timer // armed while repeats > 0, fires flushRepeats
// repeat suppression (emitLocked and the repeatSeries helpers below).
// series is the table of messages seen inside their window, keyed by
// repeatKey; seriesLRU holds the same *repeatSeries values in
// most-recently-seen-first order so eviction is O(1).
series map[string]*repeatSeries
seriesLRU *list.List
repeatTimer *time.Timer // armed at the earliest series deadline, if any
closed bool // Close ran: the timer must not write any more
// test seams
@@ -197,11 +227,13 @@ func New(stderr io.Writer, cfg Config) *Sink {
purgeLogFiles(cfg.path(), PersistPath, TmpfsPath)
}
return &Sink{
cfg: cfg,
stderr: stderr,
now: time.Now,
free: freeBytes,
window: repeatWindow,
cfg: cfg,
stderr: stderr,
series: make(map[string]*repeatSeries),
seriesLRU: list.New(),
now: time.Now,
free: freeBytes,
window: repeatWindow,
}
}
@@ -240,10 +272,9 @@ func (s *Sink) Reconfigure(cfg Config) {
if cfg == s.cfg {
return
}
// Settle any run in progress under the OLD configuration: its summary
// belongs to the destination the swallowed lines were headed for.
s.flushRepeatsLocked()
s.lastKey, s.hasLast = "", false
// Settle every series in progress under the OLD configuration: their
// summaries belong to the destination the swallowed lines were headed for.
s.flushAllSeriesLocked()
if cfg.path() != s.cfg.path() || !cfg.ToFile {
s.closeFileLocked()
}
@@ -256,10 +287,10 @@ func (s *Sink) Reconfigure(cfg Config) {
s.cfg = cfg
}
// Close flushes a pending partial line, emits the summary of a still-open run
// of repeats (the last series must never be lost) and closes the file segment.
// The sink must not be written to afterwards; a repeat timer that fires after
// Close is a no-op.
// Close flushes a pending partial line, emits the summaries of every still-open
// series of repeats (no series may be lost at shutdown) and closes the file
// segment. The sink must not be written to afterwards; a repeat timer that
// fires after Close is a no-op.
func (s *Sink) Close() error {
s.mu.Lock()
defer s.mu.Unlock()
@@ -267,7 +298,7 @@ func (s *Sink) Close() error {
s.emitLocked(s.buf)
s.buf = nil
}
s.flushRepeatsLocked()
s.flushAllSeriesLocked()
s.closed = true
if s.file != nil {
err := s.file.Close()
@@ -281,36 +312,199 @@ func (s *Sink) Close() error {
// --- internals (caller holds s.mu) -------------------------------------------
// emitLocked runs one complete line through repeat suppression and, unless it
// is swallowed as a repeat, hands it to writeLineLocked.
// is swallowed as a copy of a message already printed inside its window, hands
// it to writeLineLocked.
func (s *Sink) emitLocked(line []byte) {
if !s.cfg.ToSyslog && !s.cfg.ToFile {
return // fully off: the line is dropped, nowhere else to go
}
line = bytes.TrimSuffix(line, []byte{'\r'})
key, exempt := repeatKey(line)
switch {
case exempt:
// fatal/panic: always printed, and never becomes the head of a run —
// a dying daemon must not have its last words counted instead of said.
s.flushRepeatsLocked()
s.lastKey, s.hasLast = "", false
case s.hasLast && key == s.lastKey:
s.repeats++
if s.repeatTimer == nil {
// Arm on the 0->1 transition: the run gets a summary within one
// window even if nothing else is ever logged.
w := s.window
if w <= 0 {
w = repeatWindow
}
s.repeatTimer = time.AfterFunc(w, s.flushRepeats)
}
if exempt {
// fatal/panic: always printed, and never opens a series — a dying
// daemon must not have its last words counted instead of said. The
// pending summaries go out FIRST: the process may not live long enough
// to reach Close, and a swallowed count that is never reported is
// exactly the failure this whole mechanism exists to avoid.
s.flushAllSeriesLocked()
s.writeLineLocked(line)
return
default:
s.flushRepeatsLocked() // a different message ends the previous run
s.lastKey, s.hasLast = key, true
}
now := s.now()
if ser, ok := s.series[key]; ok {
if now.Before(ser.deadline) {
ser.count++
s.seriesLRU.MoveToFront(ser.el)
return // swallowed; the deadline did not move, so no re-arming
}
// The window elapsed without the timer having got to it (a coarse
// timer, a frozen test clock, a burst racing the callback): settle it
// here on exactly the same terms the sweep would have used.
if s.expireSeriesLocked(ser, now) {
// Still flooding: it stays suppressed, this copy opens the count of
// the new window.
ser.count = 1
s.seriesLRU.MoveToFront(ser.el)
s.armRepeatTimerLocked()
return
}
}
s.openSeriesLocked(key, now)
s.writeLineLocked(line)
s.armRepeatTimerLocked()
}
// repeatSeries is one message inside its window: the identity that is being
// collapsed, how many copies have been swallowed since the head line or the
// last summary, and when the current window ends.
type repeatSeries struct {
key string
count int
deadline time.Time
el *list.Element // this series' node in Sink.seriesLRU
}
// windowLocked is the effective suppression window (tests override s.window).
func (s *Sink) windowLocked() time.Duration {
if s.window > 0 {
return s.window
}
return repeatWindow
}
// openSeriesLocked starts a window for key, evicting the least recently seen
// series first if the table is full. An evicted series that had swallowed
// copies reports them on the way out, so the table's size limit can shorten a
// window but can never lose a count.
func (s *Sink) openSeriesLocked(key string, now time.Time) {
for len(s.series) >= maxRepeatKeys {
back := s.seriesLRU.Back()
if back == nil {
break
}
ev := back.Value.(*repeatSeries)
if ev.count > 0 {
s.summariseSeriesLocked(ev, " (repeat table full)")
}
s.dropSeriesLocked(ev)
}
ser := &repeatSeries{key: key, deadline: now.Add(s.windowLocked())}
ser.el = s.seriesLRU.PushFront(ser)
s.series[key] = ser
}
// dropSeriesLocked forgets a series entirely (its next copy prints in full).
func (s *Sink) dropSeriesLocked(ser *repeatSeries) {
s.seriesLRU.Remove(ser.el)
delete(s.series, ser.key)
}
// expireSeriesLocked settles a series whose window has ended and reports
// whether it stays open. One that swallowed copies prints their summary and
// keeps its slot for another window — an ongoing flood must be acknowledged
// every window without re-printing its head line. One that swallowed nothing is
// forgotten: the message occurs rarely enough that it needs no collapsing at
// all, and holding its slot would only push a real flood out of the table.
func (s *Sink) expireSeriesLocked(ser *repeatSeries, now time.Time) bool {
if ser.count == 0 {
s.dropSeriesLocked(ser)
return false
}
s.summariseSeriesLocked(ser, "")
ser.count = 0
ser.deadline = now.Add(s.windowLocked())
return true
}
// summariseSeriesLocked prints one summary line. It NAMES the message it counts
// (the key: the level plus the text, i.e. the line minus the uptime field that
// repeatKey drops): with several messages collapsed at once, a bare "last
// message repeated N times" would leave the reader unable to tell which line
// the number belongs to — the very confusion that made the old adjacent-only
// suppression useless in the field. note marks a summary that was forced out
// early (table full) rather than by its window.
func (s *Sink) summariseSeriesLocked(ser *repeatSeries, note string) {
unit := "times"
if ser.count == 1 {
unit = "time"
}
s.writeLineLocked([]byte(fmt.Sprintf("repeated %d %s%s: %s", ser.count, unit, note, ser.key)))
}
// onRepeatWindow is the timer callback: it settles every series whose window
// has ended, so a flood is reported while it happens instead of only when it
// stops or when the daemon closes.
func (s *Sink) onRepeatWindow() {
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
return
}
now := s.now()
for e := s.seriesLRU.Back(); e != nil; {
prev := e.Prev() // taken before a possible removal of e
ser := e.Value.(*repeatSeries)
if !now.Before(ser.deadline) {
s.expireSeriesLocked(ser, now)
}
e = prev
}
s.armRepeatTimerLocked()
}
// armRepeatTimerLocked points the single repeat timer at the earliest deadline
// in the table (and stops it when the table is empty). One timer for all series
// keeps the cost at one goroutine wake per window, not one per message.
func (s *Sink) armRepeatTimerLocked() {
if s.closed {
s.stopRepeatTimerLocked()
return
}
var earliest time.Time
for _, ser := range s.series {
if earliest.IsZero() || ser.deadline.Before(earliest) {
earliest = ser.deadline
}
}
if earliest.IsZero() {
s.stopRepeatTimerLocked()
return
}
d := earliest.Sub(s.now())
if d < 0 {
d = 0
}
if s.repeatTimer == nil {
s.repeatTimer = time.AfterFunc(d, s.onRepeatWindow)
return
}
// A callback that already started is harmless: it blocks on s.mu and then
// finds nothing expired.
s.repeatTimer.Stop()
s.repeatTimer.Reset(d)
}
func (s *Sink) stopRepeatTimerLocked() {
if s.repeatTimer != nil {
s.repeatTimer.Stop()
s.repeatTimer = nil
}
}
// flushAllSeriesLocked empties the table, printing a summary for every series
// that had swallowed copies — oldest-seen first, so the summaries come out in
// the order the messages last appeared. Used wherever the sink's state ends:
// Close, Reconfigure (the counts belong to the destination they were headed
// for) and a fatal/panic line.
func (s *Sink) flushAllSeriesLocked() {
s.stopRepeatTimerLocked()
for e := s.seriesLRU.Back(); e != nil; e = e.Prev() {
if ser := e.Value.(*repeatSeries); ser.count > 0 {
s.summariseSeriesLocked(ser, "")
}
}
s.seriesLRU.Init()
s.series = make(map[string]*repeatSeries)
}
// repeatKey reduces a raw producer line to the identity used for suppression,
@@ -327,8 +521,14 @@ func (s *Sink) emitLocked(line []byte) {
// What the key deliberately does NOT drop is the per-connection "[id duration]"
// group log.Formatter inserts for context-bound lines: those ids identify
// distinct connections, and folding them together would turn "50 connections
// failed" into one indistinguishable count. Such lines differ by id and simply
// never form a run — which is correct, they are not repeats.
// failed" into one indistinguishable count. Such lines differ by id, so each
// gets its own (single-copy, never summarised) series — which is correct, they
// are not repeats. It also means they are the messages most likely to fill the
// repeat table; that is what maxRepeatKeys and its least-recently-seen eviction
// are for.
//
// The key doubles as the text a summary names itself with, which is why it
// keeps the LEVEL and reads as a message on its own ("ERROR outbound/…: …").
//
// ANSI colour is stripped first so the key is identical for the coloured
// (terminal) and plain (procd) renderings of the same message.
@@ -357,38 +557,6 @@ func repeatKey(line []byte) (string, bool) {
return string(b), false
}
// flushRepeats is the repeat timer's callback: it closes a run that is still
// open one window after suppression started, so a flood is reported while it
// happens instead of only when it ends.
func (s *Sink) flushRepeats() {
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
return
}
s.flushRepeatsLocked()
}
// flushRepeatsLocked disarms the timer and, if lines were swallowed, prints the
// summary. lastKey is intentionally KEPT: an ongoing flood stays suppressed
// after its periodic summary instead of printing one full line per window.
func (s *Sink) flushRepeatsLocked() {
if s.repeatTimer != nil {
s.repeatTimer.Stop()
s.repeatTimer = nil
}
if s.repeats == 0 {
return
}
n := s.repeats
s.repeats = 0
unit := "times"
if n == 1 {
unit = "time"
}
s.writeLineLocked([]byte(fmt.Sprintf("last message repeated %d %s", n, unit)))
}
// writeLineLocked stamps one line with the UTC wall clock and fans it out.
func (s *Sink) writeLineLocked(line []byte) {
if !s.cfg.ToSyslog && !s.cfg.ToFile {
+252 -33
View File
@@ -425,12 +425,15 @@ func payloads(t *testing.T, body string) []string {
return out
}
var repeatSummaryRe = regexp.MustCompile(`^last message repeated (\d+) times?$`)
// A summary names the message it counts: "repeated N time(s)[ (note)]: <msg>".
var repeatSummaryRe = regexp.MustCompile(`^repeated (\d+) times?( \([^)]*\))?: (.*)$`)
// summarySum returns how many suppressed lines the summaries in payloads
// account for, and how many summary lines there were.
func summarySum(t *testing.T, lines []string) (total, count int) {
// summaryByMessage returns how many suppressed copies the summaries in payloads
// account for PER NAMED MESSAGE, and how many summary lines there were.
func summaryByMessage(t *testing.T, lines []string) (map[string]int, int) {
t.Helper()
per := make(map[string]int)
count := 0
for _, l := range lines {
m := repeatSummaryRe.FindStringSubmatch(l)
if m == nil {
@@ -440,22 +443,35 @@ func summarySum(t *testing.T, lines []string) (total, count int) {
if err != nil {
t.Fatalf("unparseable summary %q: %v", l, err)
}
total += n
per[m[3]] += n
count++
}
return per, count
}
// summarySum returns how many suppressed lines the summaries in payloads
// account for, and how many summary lines there were.
func summarySum(t *testing.T, lines []string) (total, count int) {
t.Helper()
per, count := summaryByMessage(t, lines)
for _, n := range per {
total += n
}
return total, count
}
// TestRepeatRunCollapsed: the production symptom — one broken chain repeating
// the same message — leaves ONE copy of the line plus ONE summary carrying the
// right count, in BOTH halves, and the next distinct message closes the run.
// The uptime field of the producer's prefix differs on every line: that is
// exactly what repeatKey must ignore.
// TestRepeatRunCollapsed: the simplest shape of the production symptom — one
// broken chain repeating the same message back to back — leaves ONE copy of the
// line plus ONE summary carrying the right count and naming the message, in
// BOTH halves. The uptime field of the producer's prefix differs on every line:
// that is exactly what repeatKey must ignore. An unrelated message in between
// no longer ENDS the series (a series lives for its window, not until the next
// distinct line) — it is simply printed, and the summary follows at Close.
func TestRepeatRunCollapsed(t *testing.T) {
var stderr syncBuffer
cfg := fileCfg(t, true, true, 0)
s := New(&stderr, cfg)
s.window = time.Hour // no timer flush: this test is about run boundaries
s.window = time.Hour // no timer flush: this test is about series boundaries
const msg = "outbound/urltest[chain-ewan-wg-subs-h2]: WireGuard is not ready yet"
for sec := 1; sec <= 6; sec++ {
@@ -475,8 +491,8 @@ func TestRepeatRunCollapsed(t *testing.T) {
got := payloads(t, src.body)
want := []string{
strings.TrimSuffix(engLine(1, "ERROR", msg), "\n"),
"last message repeated 5 times",
strings.TrimSuffix(engLine(7, "INFO", "something else entirely"), "\n"),
"repeated 5 times: ERROR " + msg,
}
if len(got) != len(want) {
t.Fatalf("%s: got %d lines %q, want %d %q", src.name, len(got), got, len(want), want)
@@ -489,9 +505,17 @@ func TestRepeatRunCollapsed(t *testing.T) {
}
}
// TestRepeatAlternatingNotSuppressed: A B A B is four distinct events, not a
// run — nothing may be swallowed and no summary may appear.
func TestRepeatAlternatingNotSuppressed(t *testing.T) {
// TestRepeatAlternatingSuppressedPerMessage is the reworked
// TestRepeatAlternatingNotSuppressed. Its old contract — A B A B is four
// distinct events, nothing may be swallowed — was the bug: on the router the
// flood ALWAYS alternates (one message per broken chain), so adjacency-only
// suppression collapsed nothing at all.
//
// The property that test really guarded, and that this one still guards, is
// that distinct messages are never folded into one count: A and B each get
// their OWN series, their own summary and their own N. What changed is that the
// second copy of each is now suppressed instead of printed.
func TestRepeatAlternatingSuppressedPerMessage(t *testing.T) {
var stderr syncBuffer
cfg := fileCfg(t, true, true, 0)
s := New(&stderr, cfg)
@@ -511,17 +535,171 @@ func TestRepeatAlternatingNotSuppressed(t *testing.T) {
{"stderr", stderr.String()},
} {
got := payloads(t, src.body)
if len(got) != 4 {
t.Fatalf("%s: got %d lines %q, want 4 (nothing suppressed)", src.name, len(got), got)
want := []string{
strings.TrimSuffix(engLine(1, "WARN", "alpha happened"), "\n"),
strings.TrimSuffix(engLine(2, "WARN", "beta happened"), "\n"),
"repeated 1 time: WARN alpha happened",
"repeated 1 time: WARN beta happened",
}
if _, n := summarySum(t, got); n != 0 {
t.Errorf("%s: %d summary lines for an alternating sequence: %q", src.name, n, got)
if len(got) != len(want) {
t.Fatalf("%s: got %d lines %q, want %d %q", src.name, len(got), got, len(want), want)
}
for i, want := range []string{"alpha", "beta", "alpha", "beta"} {
if !strings.Contains(got[i], want) {
t.Errorf("%s line %d = %q, want it to contain %q", src.name, i, got[i], want)
for i := range want {
if got[i] != want[i] {
t.Errorf("%s line %d = %q, want %q", src.name, i, got[i], want[i])
}
}
// The counts are per message — never merged into one "repeated 2".
per, n := summaryByMessage(t, got)
if n != 2 || per["WARN alpha happened"] != 1 || per["WARN beta happened"] != 1 {
t.Errorf("%s: summaries %v (%d lines), want one per message with N=1", src.name, per, n)
}
}
}
// TestRepeatInterleavedFloodCollapsed is the field case that motivated the
// table: the engine cycles the SAME message over three chain tags, so no two
// identical lines are ever adjacent. 300 lines must leave 3 printed heads and
// exactly 3 summaries, each naming its own message with its own count.
func TestRepeatInterleavedFloodCollapsed(t *testing.T) {
var stderr syncBuffer
cfg := fileCfg(t, true, true, 0)
s := New(&stderr, cfg)
s.window = time.Hour // the summaries come from Close, deterministically
msgs := []string{
"outbound/urltest[chain-ewan-wg-subs-h2]: WireGuard is not ready yet",
"outbound/urltest[chain-ewan-wg-subs-h3]: WireGuard is not ready yet",
"outbound/urltest[chain-ewan-wg-subs-h4]: WireGuard is not ready yet",
}
const copies = 100
for i := 0; i < copies; i++ {
for _, m := range msgs {
if _, err := s.Write([]byte(engLine(i, "ERROR", m))); err != nil {
t.Fatalf("Write: %v", err)
}
}
}
if err := s.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
for _, src := range []struct{ name, body string }{
{"file", readFile(t, cfg.Path)},
{"stderr", stderr.String()},
} {
got := payloads(t, src.body)
var want []string
for _, m := range msgs { // the heads, in first-seen order
want = append(want, strings.TrimSuffix(engLine(0, "ERROR", m), "\n"))
}
for _, m := range msgs { // the summaries, oldest-seen first
want = append(want, fmt.Sprintf("repeated %d times: ERROR %s", copies-1, m))
}
if len(got) != len(want) {
t.Fatalf("%s: got %d lines %q, want %d %q", src.name, len(got), got, len(want), want)
}
for i := range want {
if got[i] != want[i] {
t.Errorf("%s line %d = %q, want %q", src.name, i, got[i], want[i])
}
}
per, n := summaryByMessage(t, got)
if n != len(msgs) {
t.Errorf("%s: %d summary lines, want %d (one per message)", src.name, n, len(msgs))
}
for _, m := range msgs {
if per["ERROR "+m] != copies-1 {
t.Errorf("%s: message %q summarised %d copies, want %d", src.name, m, per["ERROR "+m], copies-1)
}
}
}
}
// TestRepeatTableOverflowReportsEviction: the table is bounded, and hitting the
// bound never drops a count silently — the least-recently-seen series is
// flushed with a summary that says WHY it was cut short, and the table stays at
// its limit.
func TestRepeatTableOverflowReportsEviction(t *testing.T) {
var stderr syncBuffer
cfg := fileCfg(t, true, false, 0)
s := New(&stderr, cfg)
s.window = time.Hour
// Fill the table, every key with exactly one swallowed copy to lose.
for i := 0; i < maxRepeatKeys; i++ {
msg := fmt.Sprintf("chain-%03d is down", i)
_, _ = s.Write([]byte(engLine(i, "ERROR", msg)))
_, _ = s.Write([]byte(engLine(i, "ERROR", msg)))
}
// One key too many: the oldest series (chain-000) must be evicted, and its
// swallowed copy must be reported on the way out.
_, _ = s.Write([]byte(engLine(999, "ERROR", "one key too many")))
want := "repeated 1 time (repeat table full): ERROR chain-000 is down"
got := payloads(t, stderr.String())
found := false
for _, l := range got {
if l == want {
found = true
}
}
if !found {
tail := got
if len(tail) > 4 {
tail = tail[len(tail)-4:]
}
t.Fatalf("eviction was silent: no %q (tail of the output: %q)", want, tail)
}
s.mu.Lock()
size, lru := len(s.series), s.seriesLRU.Len()
s.mu.Unlock()
if size != maxRepeatKeys || lru != maxRepeatKeys {
t.Errorf("table holds %d entries (lru %d), want the cap %d", size, lru, maxRepeatKeys)
}
_ = s.Close()
// Nothing anywhere was lost: every written line is either printed or counted.
got = payloads(t, stderr.String())
suppressed, summaries := summarySum(t, got)
printed := len(got) - summaries
if wrote := 2*maxRepeatKeys + 1; printed+suppressed != wrote {
t.Errorf("%d printed + %d suppressed = %d, want %d written", printed, suppressed, printed+suppressed, wrote)
}
}
// TestRepeatSeriesSurviveReconfigure: a live Reconfigure settles every open
// series into the destination its swallowed copies were headed for, and starts
// the next configuration with an empty table.
func TestRepeatSeriesSurviveReconfigure(t *testing.T) {
dir := t.TempDir()
cfgA := Config{ToFile: true, Path: filepath.Join(dir, "a.log")}
cfgB := Config{ToFile: true, Path: filepath.Join(dir, "b.log")}
s := New(nil, cfgA)
s.window = time.Hour
// Two interleaved series: alpha keeps 2 swallowed copies, beta keeps 1.
for sec, msg := range []string{"alpha down", "beta down", "alpha down", "beta down", "alpha down"} {
_, _ = s.Write([]byte(engLine(sec+1, "ERROR", msg)))
}
s.Reconfigure(cfgB)
_, _ = s.Write([]byte(engLine(6, "ERROR", "alpha down")))
if err := s.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
a := payloads(t, readFile(t, cfgA.Path))
perA, nA := summaryByMessage(t, a)
if nA != 2 || perA["ERROR alpha down"] != 2 || perA["ERROR beta down"] != 1 {
t.Errorf("a.log summaries %v (%d lines), want alpha=2 beta=1 flushed by Reconfigure: %q", perA, nA, a)
}
b := payloads(t, readFile(t, cfgB.Path))
if _, nB := summaryByMessage(t, b); nB != 0 {
t.Errorf("b.log carries summaries of lines written before the swap: %q", b)
}
// The new configuration starts fresh: the message prints in full again.
if len(b) != 1 || b[0] != strings.TrimSuffix(engLine(6, "ERROR", "alpha down"), "\n") {
t.Errorf("b.log = %q, want the re-printed head line only", b)
}
}
@@ -548,8 +726,8 @@ func TestRepeatLastRunSurvivesClose(t *testing.T) {
if len(got) != 2 {
t.Fatalf("%s: got %q, want the line plus one summary", src.name, got)
}
if got[1] != "last message repeated 3 times" {
t.Errorf("%s: summary = %q, want %q", src.name, got[1], "last message repeated 3 times")
if want := "repeated 3 times: ERROR dying in a loop"; got[1] != want {
t.Errorf("%s: summary = %q, want %q", src.name, got[1], want)
}
}
@@ -561,8 +739,9 @@ func TestRepeatLastRunSurvivesClose(t *testing.T) {
_, _ = s2.Write([]byte(engLine(1, "WARN", "twice only")))
_, _ = s2.Write([]byte(engLine(2, "WARN", "twice only")))
_ = s2.Close()
if got := payloads(t, stderr2.String()); len(got) != 2 || got[1] != "last message repeated 1 time" {
t.Errorf("two-line run rendered as %q, want the line plus %q", got, "last message repeated 1 time")
want := "repeated 1 time: WARN twice only"
if got := payloads(t, stderr2.String()); len(got) != 2 || got[1] != want {
t.Errorf("two-line run rendered as %q, want the line plus %q", got, want)
}
}
@@ -594,6 +773,32 @@ func TestRepeatFatalNeverSuppressed(t *testing.T) {
t.Errorf("%s: fatal/panic run produced %d summaries: %q", src.name, n, got)
}
}
// A fatal line also settles everything the table was holding BEFORE it
// speaks: the process may never reach Close, and a swallowed count that is
// never reported is the failure this mechanism exists to prevent.
var stderr2 syncBuffer
s2 := New(&stderr2, fileCfg(t, true, false, 0))
s2.window = time.Hour
for sec := 1; sec <= 3; sec++ {
_, _ = s2.Write([]byte(engLine(sec, "ERROR", "about to die")))
}
_, _ = s2.Write([]byte(engLine(4, "FATAL", "engine is gone")))
got := payloads(t, stderr2.String())
want := []string{
strings.TrimSuffix(engLine(1, "ERROR", "about to die"), "\n"),
"repeated 2 times: ERROR about to die",
strings.TrimSuffix(engLine(4, "FATAL", "engine is gone"), "\n"),
}
if len(got) != len(want) {
t.Fatalf("got %q, want %q", got, want)
}
for i := range want {
if got[i] != want[i] {
t.Errorf("line %d = %q, want %q", i, got[i], want[i])
}
}
_ = s2.Close()
}
// TestRepeatWindowFlushesOngoingRun: a flood that never stops still reports
@@ -603,28 +808,42 @@ func TestRepeatWindowFlushesOngoingRun(t *testing.T) {
var stderr syncBuffer
cfg := fileCfg(t, true, false, 0)
s := New(&stderr, cfg)
s.window = 20 * time.Millisecond
s.window = 50 * time.Millisecond
for sec := 1; sec <= 5; sec++ {
_, _ = s.Write([]byte(engLine(sec, "ERROR", "flooding")))
}
want := "repeated 4 times: ERROR flooding"
deadline := time.Now().Add(5 * time.Second)
for !strings.Contains(stderr.String(), "last message repeated") && time.Now().Before(deadline) {
for !strings.Contains(stderr.String(), "repeated ") && time.Now().Before(deadline) {
time.Sleep(5 * time.Millisecond)
}
got := payloads(t, stderr.String())
if len(got) != 2 || got[1] != "last message repeated 4 times" {
t.Fatalf("timer flush produced %q, want the line plus %q", got, "last message repeated 4 times")
if len(got) != 2 || got[1] != want {
t.Fatalf("timer flush produced %q, want the line plus %q", got, want)
}
// The run continues: still suppressed, and the next summary counts afresh.
// The flood continues. Whether the second batch lands inside the window the
// flush opened (swallowed) or after it lapsed (a fresh head line) is a race
// with the timer, so what is asserted is what must hold either way: the
// flood is reported AGAIN, every copy is accounted for exactly once, and
// every summary names this one message.
for sec := 6; sec <= 8; sec++ {
_, _ = s.Write([]byte(engLine(sec, "ERROR", "flooding")))
}
_ = s.Close()
got = payloads(t, stderr.String())
if len(got) != 3 || got[2] != "last message repeated 3 times" {
t.Fatalf("continued run produced %q, want a second summary %q", got, "last message repeated 3 times")
per, summaries := summaryByMessage(t, got)
suppressed, _ := summarySum(t, got)
printed := len(got) - summaries
if printed+suppressed != 8 {
t.Errorf("%d printed + %d suppressed = %d, want the 8 written lines: %q", printed, suppressed, printed+suppressed, got)
}
if summaries < 2 {
t.Errorf("the continuing flood was reported %d time(s), want a second summary: %q", summaries, got)
}
if len(per) != 1 || per["ERROR flooding"] != suppressed {
t.Errorf("summaries %v, want all of them naming %q", per, "ERROR flooding")
}
}
@@ -0,0 +1,380 @@
//go:build with_awg
// lx: end-to-end regression for the AmneziaWG-over-detour ClientBind fix.
//
// When an AmneziaWG endpoint runs through a detour, the WireGuard bind is our
// ClientBind (not conn.StdNetBind — that path is taken only by the DefaultDialer
// / no-detour case, see endpoint.go). Upstream ClientBind unconditionally
// stripped bytes 1-3 of every datagram (clear(b[1:4]) on receive, copy on
// Send) — those are the Cloudflare WARP "reserved" bytes. AmneziaWG 2.0 ranged
// magic headers (h1-h4) instead put a full uint32 into bytes 0-3
// (send.go RoutineEncryption: LittleEndian.PutUint32(header[0:4], magic)).
// Zeroing bytes 1-3 collapses that magic to val <= 255, which falls outside the
// configured range, so DeterminePacketTypeAndPadding returns MessageUnknownType
// and the packet is dropped — the AWG tunnel never comes up at all.
//
// This test wires two REAL wireguard-go Devices together through ClientBind on
// both ends (loopback UDP, no StdNetBind), configures ranged h1-h4 + s4 + junk
// on both, then pushes an inner IP packet through and asserts it is delivered to
// the peer's TUN within the deadline. With the pre-fix unconditional clear the
// handshake magic (h1/h2) is destroyed and delivery never happens; with the fix
// the reserved bytes are left alone for a plain-AWG (reserved == [0,0,0])
// endpoint and the packet arrives.
package wireguard
import (
"context"
"crypto/rand"
"encoding/hex"
"fmt"
"net"
"net/netip"
"os"
"testing"
"time"
"github.com/sagernet/sing/common/logger"
M "github.com/sagernet/sing/common/metadata"
N "github.com/sagernet/sing/common/network"
"github.com/sagernet/sing/service/pause"
"github.com/sagernet/wireguard-go/device"
"github.com/sagernet/wireguard-go/tun"
"golang.org/x/crypto/curve25519"
)
// ---------------------------------------------------------------------------
// loopback UDP dialer: replaces N.Dialer so ClientBind.connect() gets a real
// UDP socket bound to 127.0.0.1, connected to the peer's loopback listener.
// This is the "detour" stand-in — the point is only that the bind is ClientBind
// and NOT conn.StdNetBind.
// ---------------------------------------------------------------------------
// loopbackDialer binds ClientBind's listen socket to a FIXED loopback port so
// the two binds can address each other by a known port. We use the non-connect
// ClientBind path (isConnect == false), which routes outbound via
// PacketConn.WriteTo to the wireguard-configured endpoint addr:port — matching
// each bind's fixed listen port. (The connect path dials from an ephemeral
// source port, so neither side would ever bind the reserved ports.)
type loopbackDialer struct {
listenPort int
}
func (loopbackDialer) DialContext(ctx context.Context, network string, destination M.Socksaddr) (net.Conn, error) {
var d net.Dialer
d.LocalAddr = &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)}
return d.DialContext(ctx, "udp", destination.String())
}
func (d loopbackDialer) ListenPacket(ctx context.Context, destination M.Socksaddr) (net.PacketConn, error) {
var lc net.ListenConfig
return lc.ListenPacket(ctx, "udp", fmt.Sprintf("127.0.0.1:%d", d.listenPort))
}
var _ N.Dialer = loopbackDialer{}
// ---------------------------------------------------------------------------
// channelTUN: a minimal tun.Device. Read() blocks handing out inbound IP
// packets (packets we inject to be encrypted and sent to the peer); Write()
// captures decrypted inner packets the Device delivers — that is the delivery
// signal the test waits on. A leading `offset` region is reserved exactly like
// a real TUN, which the Device uses for its own headers.
// ---------------------------------------------------------------------------
type channelTUN struct {
name string
mtu int
inbound chan []byte // packets to hand to the Device via Read (to encrypt+send)
outbound chan []byte // packets the Device delivered via Write (decrypted inner)
events chan tun.Event
closed chan struct{}
}
func newChannelTUN(name string, mtu int) *channelTUN {
t := &channelTUN{
name: name,
mtu: mtu,
inbound: make(chan []byte, 16),
outbound: make(chan []byte, 16),
events: make(chan tun.Event, 4),
closed: make(chan struct{}),
}
t.events <- tun.EventUp
return t
}
func (t *channelTUN) File() *os.File { return nil }
func (t *channelTUN) Read(bufs [][]byte, sizes []int, offset int) (int, error) {
select {
case <-t.closed:
return 0, os.ErrClosed
case pkt := <-t.inbound:
n := copy(bufs[0][offset:], pkt)
sizes[0] = n
return 1, nil
}
}
func (t *channelTUN) Write(bufs [][]byte, offset int) (int, error) {
for _, b := range bufs {
if len(b) <= offset {
continue
}
pkt := make([]byte, len(b)-offset)
copy(pkt, b[offset:])
select {
case t.outbound <- pkt:
case <-t.closed:
return 0, os.ErrClosed
default:
}
}
return len(bufs), nil
}
func (t *channelTUN) MTU() (int, error) { return t.mtu, nil }
func (t *channelTUN) Name() (string, error) { return t.name, nil }
func (t *channelTUN) Events() <-chan tun.Event { return t.events }
func (t *channelTUN) BatchSize() int { return 1 }
func (t *channelTUN) Close() error {
select {
case <-t.closed:
default:
close(t.closed)
close(t.events)
}
return nil
}
var _ tun.Device = (*channelTUN)(nil)
// ---------------------------------------------------------------------------
// key material
// ---------------------------------------------------------------------------
type wgKeypair struct {
private [32]byte
public [32]byte
}
func genKeypair(t *testing.T) wgKeypair {
t.Helper()
var kp wgKeypair
if _, err := rand.Read(kp.private[:]); err != nil {
t.Fatalf("read random: %v", err)
}
// curve25519 clamping (as WireGuard does for private keys).
kp.private[0] &= 248
kp.private[31] &= 127
kp.private[31] |= 64
pub, err := curve25519.X25519(kp.private[:], curve25519.Basepoint)
if err != nil {
t.Fatalf("derive public key: %v", err)
}
copy(kp.public[:], pub)
return kp
}
// awgObfLines are the AmneziaWG 2.0 obfuscation knobs, identical on both ends so
// the handshake magic ranges line up. h1-h4 are ranged (AWG 2.0) so the magic
// occupies the full uint32 in bytes 0-3 — exactly the bytes the buggy
// ClientBind used to zero.
const awgObfLines = "" +
"\njc=4" +
"\njmin=8" +
"\njmax=80" +
"\ns4=12" +
"\nh1=1888111000-1888111100" +
"\nh2=1888122000-1888122100" +
"\nh3=1888133000-1888133100" +
"\nh4=1888222333-1888222444"
// buildDevice constructs a real wireguard-go Device driven by ClientBind (the
// detour path). It listens on a loopback UDP port and dials the peer's port.
func buildDevice(t *testing.T, ctx context.Context, name string, tunDev tun.Device, self wgKeypair, peerPub [32]byte, listenPort int, peerAddr netip.AddrPort) (*device.Device, *ClientBind) {
t.Helper()
// ClientBind (the detour bind), NOT conn.StdNetBind. Non-connect mode so the
// bind listens on a fixed loopback port and both ends can reach each other.
// reserved stays [0,0,0]: this is plain AmneziaWG, NOT WARP — exactly the case
// the pre-fix unconditional clear(b[1:4]) corrupted.
bind := NewClientBind(ctx, logger.NOP(), loopbackDialer{listenPort: listenPort}, false, peerAddr, [3]uint8{})
dev := device.NewDevice(ctx, tunDev, bind, &device.Logger{
Verbosef: func(string, ...any) {},
Errorf: func(format string, args ...any) { t.Logf("["+name+"] "+format, args...) },
}, 0)
ipc := "private_key=" + hex.EncodeToString(self.private[:]) +
fmt.Sprintf("\nlisten_port=%d", listenPort) +
awgObfLines +
"\npublic_key=" + hex.EncodeToString(peerPub[:]) +
"\nendpoint=" + peerAddr.String() +
"\npersistent_keepalive_interval=1" +
"\nallowed_ip=0.0.0.0/0"
if err := dev.IpcSet(ipc); err != nil {
t.Fatalf("[%s] IpcSet: %v", name, err)
}
if err := dev.Up(); err != nil {
t.Fatalf("[%s] device up: %v", name, err)
}
return dev, bind
}
// TestAwgDetourClientBindDelivers is the red/green e2e. It stands up two AWG
// Devices linked through ClientBind (loopback UDP), pushes an inner IP packet,
// and asserts delivery within 20s. See the file header for why the pre-fix
// unconditional clear(b[1:4]) makes this impossible.
func TestAwgDetourClientBindDelivers(t *testing.T) {
// A pause manager must be in the context: ClientBind.receive/Send and the
// Device both pull it via service.FromContext and call WaitActive().
ctx := pause.ContextWithDefaultManager(context.Background())
// Two loopback UDP listeners just to reserve ports; ClientBind opens its own
// sockets via the dialer, so we only need the port numbers to be free and
// wire each side to the other's port.
portA := reserveLoopbackUDPPort(t)
portB := reserveLoopbackUDPPort(t)
addrA := netip.AddrPortFrom(netip.MustParseAddr("127.0.0.1"), uint16(portA))
addrB := netip.AddrPortFrom(netip.MustParseAddr("127.0.0.1"), uint16(portB))
kpA := genKeypair(t)
kpB := genKeypair(t)
const mtu = 1420
tunA := newChannelTUN("wgA", mtu)
tunB := newChannelTUN("wgB", mtu)
// A: 10.0.0.1, B: 10.0.0.2 (allowed_ip 0.0.0.0/0 on both, so routing is trivial).
devA, _ := buildDevice(t, ctx, "A", tunA, kpA, kpB.public, portA, addrB)
devB, _ := buildDevice(t, ctx, "B", tunB, kpB, kpA.public, portB, addrA)
defer devA.Close()
defer devB.Close()
defer tunA.Close()
defer tunB.Close()
// Craft an inner IPv4/UDP packet from 10.0.0.1 -> 10.0.0.2 carrying a marker.
marker := []byte("LX-AWG-DETOUR-E2E")
pkt := buildIPv4UDP(
netip.MustParseAddr("10.0.0.1"), netip.MustParseAddr("10.0.0.2"),
4711, 4712, marker,
)
// Feed A's TUN so the Device encrypts it and sends it (over ClientBind) to B.
// Resend periodically: the first datagrams race the handshake, and until the
// handshake completes there is no keypair to encrypt transport data.
sendDone := make(chan struct{})
go func() {
ticker := time.NewTicker(200 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-sendDone:
return
default:
}
buf := make([]byte, len(pkt))
copy(buf, pkt)
select {
case tunA.inbound <- buf:
case <-sendDone:
return
}
select {
case <-ticker.C:
case <-sendDone:
return
}
}
}()
defer close(sendDone)
deadline := time.After(20 * time.Second)
for {
select {
case got := <-tunB.outbound:
if containsMarker(got, marker) {
return // GREEN: inner packet delivered end-to-end through ClientBind
}
// Ignore non-marker traffic (keepalives never reach TUN, but be safe).
case <-deadline:
t.Fatal("timeout: inner AWG packet was not delivered to peer TUN within 20s " +
"(pre-fix ClientBind zeroes bytes 1-3, collapsing the ranged h1-h4 magic " +
"below its range so the handshake/transport packets are dropped)")
}
}
}
// reserveLoopbackUDPPort grabs a free UDP port on loopback and releases it, so
// ClientBind's own listen socket can claim it. There is a tiny race window, but
// on loopback in a test it is not a practical problem.
func reserveLoopbackUDPPort(t *testing.T) int {
t.Helper()
c, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 0})
if err != nil {
t.Fatalf("reserve udp port: %v", err)
}
port := c.LocalAddr().(*net.UDPAddr).Port
_ = c.Close()
return port
}
// buildIPv4UDP assembles a minimal IPv4 + UDP datagram (with checksums) carrying
// payload, so a real WireGuard Device routes and delivers it as a valid IP packet.
func buildIPv4UDP(src, dst netip.Addr, srcPort, dstPort uint16, payload []byte) []byte {
udpLen := 8 + len(payload)
totalLen := 20 + udpLen
b := make([]byte, totalLen)
// IPv4 header
b[0] = 0x45 // version 4, IHL 5
b[1] = 0x00
putU16(b[2:], uint16(totalLen))
putU16(b[4:], 0) // id
putU16(b[6:], 0) // flags/frag
b[8] = 64 // TTL
b[9] = 17 // protocol UDP
putU16(b[10:], 0) // checksum (fill below)
copy(b[12:16], src.AsSlice())
copy(b[16:20], dst.AsSlice())
putU16(b[10:], ipChecksum(b[:20]))
// UDP header
putU16(b[20:], srcPort)
putU16(b[22:], dstPort)
putU16(b[24:], uint16(udpLen))
putU16(b[26:], 0) // checksum optional for IPv4, leave 0
copy(b[28:], payload)
return b
}
func putU16(b []byte, v uint16) {
b[0] = byte(v >> 8)
b[1] = byte(v)
}
func ipChecksum(h []byte) uint16 {
var sum uint32
for i := 0; i+1 < len(h); i += 2 {
sum += uint32(h[i])<<8 | uint32(h[i+1])
}
for sum>>16 != 0 {
sum = (sum & 0xffff) + (sum >> 16)
}
return ^uint16(sum)
}
func containsMarker(pkt, marker []byte) bool {
if len(pkt) < len(marker) {
return false
}
for i := 0; i+len(marker) <= len(pkt); i++ {
if string(pkt[i:i+len(marker)]) == string(marker) {
return true
}
}
return false
}
+61 -2
View File
@@ -50,6 +50,53 @@ func NewClientBind(ctx context.Context, logger logger.Logger, dialer N.Dialer, i
}
}
// hasReserved reports whether any Cloudflare "reserved" value is set. The
// receive path must only zero bytes 1-3 when a reserved value exists (WARP);
// otherwise an AmneziaWG magic header that lands in bytes 1-3 (small s1/s2/s4
// padding) would be corrupted and the packet dropped. The send path stamps the
// bytes under the same condition.
//
// This is the twin of StdNetBind.hasReserved (see the vendored
// submodules/wireguard-go/conn/bind_std.go, hasReserved + its callers in
// receiveIP and Send). The two binds implement ONE contract — "never touch
// bytes 1-3 unless WARP is configured" — and must be read, and fixed, as a
// pair. They were not: bind_std.go got the gate first and ClientBind kept the
// unconditional clear, which is exactly how the bug below survived.
//
// Why the bug is detour-only. transport/wireguard/endpoint.go (Endpoint.Start)
// picks the bind with common.Cast[dialer.WireGuardListener](e.options.Dialer):
//
// - no detour -> the dialer is *dialer.DefaultDialer, the only type with a
// WireGuardControl method, so the cast succeeds -> conn.StdNetBind (gated,
// healthy);
// - detour set -> dialer.NewWithOptions builds a *dialer.DetourDialer, whose
// Upstream() is the detour outbound (e.g. *direct.Outbound). Neither
// implements WireGuardControl, so common.Cast walks the upstream chain and
// fails -> ClientBind (this file). An AmneziaWG node therefore works
// standalone and dies the moment it is placed behind an egress/chain hop.
//
// Why handshakes survived and transport did not. AmneziaWG writes the magic as
// a little-endian uint32 at packet[padding:], where padding is s1 (initiation),
// s2 (response) or s4 (transport) — see device.DeterminePacketTypeAndPadding.
// With the usual s1/s2 > 0 the magic sits well past byte 3, so clearing bytes
// 1-3 only scribbles on the random junk prefix and the handshake completes. s4
// is 0 by default, so a transport packet puts the magic in bytes 0-3: zeroing
// 1-3 collapses it to its low byte (<= 255), which falls outside every h4
// range, the peer classifies it MessageUnknownType and drops it silently. The
// session looks established locally, every data packet vanishes, and the peer
// re-handshakes forever.
func (c *ClientBind) hasReserved() bool {
if c.reserved != [3]uint8{} {
return true
}
for _, reserved := range c.reservedForEndpoint {
if reserved != [3]uint8{} {
return true
}
}
return false
}
func (c *ClientBind) connect() (*wireConn, error) {
serverConn := c.conn
if serverConn != nil {
@@ -134,7 +181,12 @@ func (c *ClientBind) receive(packets [][]byte, sizes []int, eps []conn.Endpoint)
return
}
sizes[0] = n
if n > 3 {
// lx: only strip the Cloudflare "reserved" bytes when a reserved value is
// actually configured (WARP). AmneziaWG writes a full uint32 magic header
// into bytes 0-3 of a transport packet (s4 defaults to 0); unconditionally
// clearing 1-3, as upstream did, destroys it and the endpoint drops every
// inbound data packet. StdNetBind gates the same clear on hasReserved().
if n > 3 && c.hasReserved() {
b := packets[0]
clear(b[1:4])
}
@@ -179,7 +231,14 @@ func (c *ClientBind) Send(bufs [][]byte, ep conn.Endpoint, offset int) error {
if !loaded {
reserved = c.reserved
}
copy(buf[1:4], reserved[:])
// lx: only stamp the reserved bytes when non-zero (WARP). For a
// plain WG / AmneziaWG endpoint reserved is [0,0,0]; overwriting
// bytes 1-3 zeroes the upper bytes of the AWG magic header and the
// peer drops the packet. See the matching guard in receive() and
// the hasReserved() comment for why this bites only on detour.
if reserved != [3]uint8{} {
copy(buf[1:4], reserved[:])
}
}
_, err = udpConn.WriteToUDPAddrPort(buf, destination)
if err != nil {
@@ -0,0 +1,260 @@
// lx: guards the AmneziaWG-vs-reserved fix in ClientBind — bytes 1-3 (the
// Cloudflare "reserved" field) must only be touched when a reserved value is
// configured (WARP). For a plain WG / AmneziaWG endpoint they carry the upper
// bytes of a magic header, and clearing them breaks the tunnel.
//
// These tests drive the REAL ClientBind.Send / ClientBind.receive over loopback
// UDP, so they observe what actually lands on the wire rather than restating
// the condition. Before the hasReserved() gate they are red:
//
// Send: magic corrupted on the wire: got 40 (0x00000028), want 1618116904
// receive: magic corrupted on receive: got 40 (0x00000028), want 1618116904
//
// 0x28 is the low byte of the magic — exactly the failure seen on the router,
// where a node with h4=1618116899-1618116949 put 56 (0x38, inside the range's
// low-byte window 0x23..0x55) on the wire and the peer dropped every packet.
package wireguard
import (
"context"
"encoding/binary"
"net"
"net/netip"
"os"
"testing"
"time"
"github.com/sagernet/sing/common/logger"
M "github.com/sagernet/sing/common/metadata"
N "github.com/sagernet/sing/common/network"
"github.com/sagernet/wireguard-go/conn"
)
// awgTransportMagic is a value from a real AmneziaWG h4 range
// (1618116899-1618116949 = 0x60728123..0x60728155). It uses bytes 1-3, so any
// unconditional reserved-clear collapses it to 0x28 = 40.
const awgTransportMagic uint32 = 1618116904 // 0x60728128
func testAddrPort() netip.AddrPort {
return netip.MustParseAddrPort("192.0.2.1:51820")
}
func newTestClientBind(reserved [3]uint8) *ClientBind {
return &ClientBind{
reservedForEndpoint: make(map[netip.AddrPort][3]uint8),
reserved: reserved,
}
}
// testLoopbackDialer is the stand-in for a detour dialer: the only thing that
// matters is that it is NOT a dialer.DefaultDialer, so Endpoint.Start would
// pick ClientBind over conn.StdNetBind. It hands out a real loopback UDP
// socket and remembers it so the test can address the bind.
type testLoopbackDialer struct {
packetConn net.PacketConn
}
func (d *testLoopbackDialer) DialContext(ctx context.Context, network string, destination M.Socksaddr) (net.Conn, error) {
return nil, os.ErrInvalid // the tests use the non-connect (ListenPacket) path
}
func (d *testLoopbackDialer) ListenPacket(ctx context.Context, destination M.Socksaddr) (net.PacketConn, error) {
var lc net.ListenConfig
packetConn, err := lc.ListenPacket(ctx, "udp", "127.0.0.1:0")
if err != nil {
return nil, err
}
d.packetConn = packetConn
return packetConn, nil
}
var _ N.Dialer = (*testLoopbackDialer)(nil)
// newLoopbackClientBind builds an Open()'d ClientBind over loopback UDP. The
// struct is populated directly (rather than via NewClientBind) so the test does
// not need a service registry for the pause manager; every field connect(),
// Send() and receive() touch is set.
func newLoopbackClientBind(t *testing.T, reserved [3]uint8) (*ClientBind, *testLoopbackDialer) {
t.Helper()
dialer := &testLoopbackDialer{}
bind := &ClientBind{
ctx: context.Background(),
logger: logger.NOP(),
dialer: dialer,
reservedForEndpoint: make(map[netip.AddrPort][3]uint8),
done: make(chan struct{}),
reserved: reserved,
}
if _, _, err := bind.Open(0); err != nil {
t.Fatalf("Open: %v", err)
}
t.Cleanup(func() { bind.Close() })
if _, err := bind.connect(); err != nil {
t.Fatalf("connect: %v", err)
}
return bind, dialer
}
func awgTransportPacket() []byte {
// 32 bytes: type(4) + receiver(4) + counter(8) + poly1305(16) — a keepalive,
// the exact shape that failed on the router.
packet := make([]byte, 32)
binary.LittleEndian.PutUint32(packet[0:4], awgTransportMagic)
return packet
}
func TestClientBindHasReserved(t *testing.T) {
if newTestClientBind([3]uint8{}).hasReserved() {
t.Fatal("empty reserved must report false (AmneziaWG / plain WG)")
}
if !newTestClientBind([3]uint8{0, 0, 1}).hasReserved() {
t.Fatal("non-zero global reserved must report true (WARP)")
}
bind := newTestClientBind([3]uint8{})
bind.reservedForEndpoint[testAddrPort()] = [3]uint8{9, 9, 9}
if !bind.hasReserved() {
t.Fatal("non-zero per-endpoint reserved must report true (WARP)")
}
}
// TestClientBindSendPreservesAWGMagic is the send-side regression: with no
// reserved configured, the 4-byte AmneziaWG magic header of a transport packet
// must reach the wire byte-for-byte.
func TestClientBindSendPreservesAWGMagic(t *testing.T) {
peer, err := net.ListenPacket("udp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen peer: %v", err)
}
defer peer.Close()
peerAddr, err := netip.ParseAddrPort(peer.LocalAddr().String())
if err != nil {
t.Fatalf("parse peer addr: %v", err)
}
bind, _ := newLoopbackClientBind(t, [3]uint8{})
if err = bind.Send([][]byte{awgTransportPacket()}, remoteEndpoint(peerAddr), 0); err != nil {
t.Fatalf("Send: %v", err)
}
buffer := make([]byte, 128)
if err = peer.SetReadDeadline(time.Now().Add(10 * time.Second)); err != nil {
t.Fatalf("set deadline: %v", err)
}
n, _, err := peer.ReadFrom(buffer)
if err != nil {
t.Fatalf("read from peer: %v", err)
}
if n != 32 {
t.Fatalf("unexpected datagram size on the wire: got %d, want 32", n)
}
if got := binary.LittleEndian.Uint32(buffer[0:4]); got != awgTransportMagic {
t.Fatalf("magic corrupted on the wire: got %d (0x%08x), want %d (0x%08x)",
got, got, awgTransportMagic, awgTransportMagic)
}
}
// TestClientBindSendStampsReservedForWARP pins the unchanged WARP behaviour:
// a non-zero reserved value is still written into bytes 1-3 on send.
func TestClientBindSendStampsReservedForWARP(t *testing.T) {
peer, err := net.ListenPacket("udp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen peer: %v", err)
}
defer peer.Close()
peerAddr, err := netip.ParseAddrPort(peer.LocalAddr().String())
if err != nil {
t.Fatalf("parse peer addr: %v", err)
}
bind, _ := newLoopbackClientBind(t, [3]uint8{1, 2, 3})
if err = bind.Send([][]byte{awgTransportPacket()}, remoteEndpoint(peerAddr), 0); err != nil {
t.Fatalf("Send: %v", err)
}
buffer := make([]byte, 128)
if err = peer.SetReadDeadline(time.Now().Add(10 * time.Second)); err != nil {
t.Fatalf("set deadline: %v", err)
}
if _, _, err = peer.ReadFrom(buffer); err != nil {
t.Fatalf("read from peer: %v", err)
}
want := [3]uint8{1, 2, 3}
if got := [3]uint8{buffer[1], buffer[2], buffer[3]}; got != want {
t.Fatalf("WARP reserved not stamped: got %v, want %v", got, want)
}
if wantByte0 := byte(awgTransportMagic & 0xFF); buffer[0] != wantByte0 {
t.Fatalf("byte 0 must be untouched: got 0x%02x, want 0x%02x", buffer[0], wantByte0)
}
}
// TestClientBindReceivePreservesAWGMagic is the receive-side regression: an
// inbound AmneziaWG transport packet must reach the device with its magic
// intact when no reserved value is configured.
func TestClientBindReceivePreservesAWGMagic(t *testing.T) {
bind, dialer := newLoopbackClientBind(t, [3]uint8{})
bindAddr, err := net.ResolveUDPAddr("udp", dialer.packetConn.LocalAddr().String())
if err != nil {
t.Fatalf("resolve bind addr: %v", err)
}
sender, err := net.ListenPacket("udp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen sender: %v", err)
}
defer sender.Close()
if _, err = sender.WriteTo(awgTransportPacket(), bindAddr); err != nil {
t.Fatalf("write to bind: %v", err)
}
// Bound the blocking ReadFrom inside receive() so a lost datagram fails the
// test instead of hanging it.
if err = dialer.packetConn.SetReadDeadline(time.Now().Add(10 * time.Second)); err != nil {
t.Fatalf("set deadline: %v", err)
}
packets := [][]byte{make([]byte, 2048)}
sizes := make([]int, 1)
endpoints := make([]conn.Endpoint, 1)
count, err := bind.receive(packets, sizes, endpoints)
if err != nil {
t.Fatalf("receive: %v", err)
}
if count != 1 || sizes[0] != 32 {
t.Fatalf("unexpected receive result: count=%d size=%d", count, sizes[0])
}
if got := binary.LittleEndian.Uint32(packets[0][0:4]); got != awgTransportMagic {
t.Fatalf("magic corrupted on receive: got %d (0x%08x), want %d (0x%08x)",
got, got, awgTransportMagic, awgTransportMagic)
}
}
// TestClientBindReceiveStripsReservedForWARP pins the unchanged WARP behaviour
// on the receive side: with a reserved value configured, bytes 1-3 are cleared.
func TestClientBindReceiveStripsReservedForWARP(t *testing.T) {
bind, dialer := newLoopbackClientBind(t, [3]uint8{0, 0, 1})
bindAddr, err := net.ResolveUDPAddr("udp", dialer.packetConn.LocalAddr().String())
if err != nil {
t.Fatalf("resolve bind addr: %v", err)
}
sender, err := net.ListenPacket("udp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen sender: %v", err)
}
defer sender.Close()
if _, err = sender.WriteTo(awgTransportPacket(), bindAddr); err != nil {
t.Fatalf("write to bind: %v", err)
}
if err = dialer.packetConn.SetReadDeadline(time.Now().Add(10 * time.Second)); err != nil {
t.Fatalf("set deadline: %v", err)
}
packets := [][]byte{make([]byte, 2048)}
sizes := make([]int, 1)
endpoints := make([]conn.Endpoint, 1)
if _, err = bind.receive(packets, sizes, endpoints); err != nil {
t.Fatalf("receive: %v", err)
}
if got := binary.LittleEndian.Uint32(packets[0][0:4]); got != awgTransportMagic&0xFF {
t.Fatalf("WARP reserved not stripped: got %d, want %d", got, awgTransportMagic&0xFF)
}
}