Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bcc9df9282 | ||
|
|
ee3641fe45 | ||
|
|
439f62238f | ||
|
|
d971eb85ee |
+19
-7
@@ -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).
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
@@ -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
@@ -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
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user