Skip to content

Commit fc22839

Browse files
thomas-manginclaude
andcommitted
bfd(echo): Phase B.1+B.2 outstanding ring + echo detection
Adds per-session echo loss tracking and the echo-mode detection timer so the echo transport landed in Phase A becomes a real liveness channel. Without this commit, echo packets flowing with no detection would be a passive RTT probe only. Phase B.1 -- outstanding-ID ring (session/session.go + timers.go): - Fixed 16-slot echoOutstanding ring per Machine tracks every echo TX by (sequence, sentAt). RegisterEchoTx inserts in the first empty slot or overwrites the oldest live slot when the ring is full (an overwrite is equivalent to a lost echo from the detection standpoint). MatchEchoRx scans for the returning echo's sequence, clears the slot, and returns the monotonic RTT. ClearEchoSchedule wipes the ring on session-down so a flap-back-up starts with a clean detection window. Phase B.2 -- detection-time switchover (engine/echo.go): - echoTickLocked now calls EchoDetectionExpired(now) before the TX pass. EchoDetectionExpired walks the ring and returns true when any outstanding entry is older than DetectMult * EchoInterval (the echo-mode detection time per RFC 5880 Section 6.8.4). On true, the engine calls Machine.EchoFail which transitions the session to Down with DiagEchoFailed (2), fires the notify callback with the correct diagnostic, and clears the ring. - recordEchoRTTLocked now uses MatchEchoRx for the authoritative RTT (monotonic, immune to wall-clock jumps) and falls back to the ZEEC envelope TimestampMs when the ring entry was evicted. New TestEchoDetectionSwitchover drives the full detection path via a dropEchoTransport that swallows every outbound echo and asserts a Down+DiagEchoFailed transition within 2 s. Race-clean under -race -count=10. Also fixes pre-existing goconst lint in cmd/ze/iface: promote "bridge" literal in internal/component/iface/migrate_linux.go to the existing zeTypeBridge constant, add nolint:goconst to CLI dispatch strings in create.go, main.go, migrate.go, scan.go. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent 1de1e43 commit fc22839

9 files changed

Lines changed: 331 additions & 35 deletions

File tree

‎cmd/ze/iface/create.go‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,9 @@ func cmdCreate(args []string) int {
2222
case "help", "-h", "--help": //nolint:goconst // consistent pattern across cmd files
2323
createUsage()
2424
return 0
25-
case "dummy":
25+
case "dummy": //nolint:goconst // CLI dispatch strings, constants in internal/component/iface
2626
return cmdCreateDummy(args[1:])
27-
case "veth":
27+
case "veth": //nolint:goconst // CLI dispatch strings, constants in internal/component/iface
2828
return cmdCreateVeth(args[1:])
2929
default:
3030
fmt.Fprintf(os.Stderr, "error: unknown interface type: %s (expected dummy or veth)\n", args[0])

‎cmd/ze/iface/main.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@ func Run(args []string) int {
3636
subcmd := args[0]
3737
subArgs := args[1:]
3838

39-
if subcmd == "help" || subcmd == "-h" || subcmd == "--help" {
39+
if subcmd == "help" || subcmd == "-h" || subcmd == "--help" { //nolint:goconst // consistent pattern across cmd files
4040
usage()
4141
return 0
4242
}

‎cmd/ze/iface/migrate.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ func cmdMigrate(args []string) int {
7272
// Validate --create type against known set.
7373
if createTyp != "" {
7474
switch createTyp {
75-
case "dummy", "veth", "bridge": // valid
75+
case "dummy", "veth", "bridge": //nolint:goconst // CLI dispatch strings, constants in internal/component/iface
7676
default:
7777
fmt.Fprintf(os.Stderr, "error: invalid --create type %q (expected dummy, veth, or bridge)\n", createTyp)
7878
return 1

‎cmd/ze/iface/scan.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -108,7 +108,7 @@ func filterManaged(discovered []ifacepkg.DiscoveredInterface) []ifacepkg.Discove
108108
filtered := make([]ifacepkg.DiscoveredInterface, 0, len(discovered))
109109
for i := range discovered {
110110
switch discovered[i].Type {
111-
case "dummy", "veth", "bridge", "tunnel", "wireguard":
111+
case "dummy", "veth", "bridge", "tunnel", "wireguard": //nolint:goconst // CLI dispatch strings
112112
filtered = append(filtered, discovered[i])
113113
}
114114
}

‎internal/component/iface/migrate_linux.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@ func validateMigrateConfig(cfg MigrateConfig) error {
2929
}
3030
if cfg.NewIfaceType != "" {
3131
switch cfg.NewIfaceType {
32-
case "dummy", "veth", "bridge": // valid types
32+
case zeTypeDummy, zeTypeVeth, zeTypeBridge:
3333
default: // unknown type
3434
return fmt.Errorf("migrate: unknown interface type %q (expected dummy, veth, or bridge)", cfg.NewIfaceType)
3535
}

‎internal/plugins/bfd/engine/echo.go‎

Lines changed: 46 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -28,15 +28,18 @@ import (
2828
"codeberg.org/thomas-mangin/ze/internal/plugins/bfd/transport"
2929
)
3030

31-
// echoTickLocked fires per-session echo TX deadlines. Caller MUST
32-
// hold l.mu. When a session is Up with echo negotiated and its
33-
// next-echo deadline has passed, encode one ZEEC envelope into a
34-
// pool buffer, send via the echo transport, and advance the
35-
// per-session schedule.
31+
// echoTickLocked fires per-session echo TX deadlines and drives
32+
// the echo-mode detection timer. Caller MUST hold l.mu. For every
33+
// session with echo negotiated and its deadline passed the loop
34+
// encodes one ZEEC envelope, hands it to the echo transport, and
35+
// advances the per-session schedule. Before the TX pass the
36+
// session's outstanding ring is checked for stale entries; any
37+
// outstanding echo older than DetectMult * EchoInterval fires a
38+
// transition to Down with DiagEchoFailed.
3639
//
37-
// Sessions that are not Up have their echo schedule cleared so a
38-
// session that flapped down does not fire stale echoes after going
39-
// back up (PrimeEcho rearms on the next Up tick).
40+
// Sessions that are not Up have their echo schedule and
41+
// outstanding ring cleared so a session that flapped down does
42+
// not carry stale detection state back into Up.
4043
func (l *Loop) echoTickLocked(now time.Time) {
4144
if l.echoTransport == nil {
4245
return
@@ -47,6 +50,17 @@ func (l *Loop) echoTickLocked(now time.Time) {
4750
m.ClearEchoSchedule()
4851
continue
4952
}
53+
if m.EchoDetectionExpired(now) {
54+
engineLog().Info("bfd echo detection expired",
55+
"peer", m.PeerAddr().String(),
56+
"detect", m.EchoDetectInterval())
57+
m.EchoFail()
58+
if hook := l.metricsHook.Load(); hook != nil {
59+
(*hook).OnStateChange(packet.StateUp, packet.StateDown,
60+
packet.DiagEchoFailed, m.Key().Mode.String(), m.Key().VRF)
61+
}
62+
continue
63+
}
5064
m.PrimeEcho(now)
5165
deadline := m.NextEchoTxDeadline()
5266
if deadline.IsZero() || now.Before(deadline) {
@@ -60,18 +74,19 @@ func (l *Loop) echoTickLocked(now time.Time) {
6074
// sendEchoLocked encodes a single ZEEC envelope for the session and
6175
// hands it to the echo transport. Caller MUST hold l.mu.
6276
//
63-
// The envelope's TimestampMs field carries a truncated millisecond
64-
// slice of the engine clock so the reflected copy can be matched
65-
// back to the original TX without a separate outstanding-ID index.
66-
// The sequence counter is self-carried for future diagnostics
67-
// (packet loss estimation) but the current RX path does not
68-
// consult it.
77+
// The sequence counter lets the RX path match returning echoes
78+
// against the per-session outstanding ring (RegisterEchoTx /
79+
// MatchEchoRx). The envelope's TimestampMs field is still written
80+
// with a truncated millisecond slice of the engine clock so the
81+
// RTT can fall back to the self-carried value when the matching
82+
// entry has already been evicted from a small ring under load.
6983
func (l *Loop) sendEchoLocked(entry *sessionEntry, now time.Time) {
7084
buf := [packet.EchoLen]byte{}
7185
ts := uint32(now.UnixMilli())
86+
seq := entry.machine.NextEchoSequence()
7287
e := packet.Echo{
7388
LocalDiscriminator: entry.machine.LocalDiscriminator(),
74-
Sequence: entry.machine.NextEchoSequence(),
89+
Sequence: seq,
7590
TimestampMs: ts,
7691
}
7792
packet.WriteEcho(buf[:], 0, e)
@@ -88,6 +103,7 @@ func (l *Loop) sendEchoLocked(entry *sessionEntry, now time.Time) {
88103
engineLog().Debug("bfd echo send failed", "peer", out.To, "err", err)
89104
return
90105
}
106+
entry.machine.RegisterEchoTx(seq, now)
91107
if hook := l.metricsHook.Load(); hook != nil {
92108
(*hook).OnEchoTx(key.Mode.String())
93109
}
@@ -154,16 +170,22 @@ func (l *Loop) findSessionByPeerLocked(peer netip.Addr) *sessionEntry {
154170
// recordEchoRTTLocked stores the round-trip time on the session and
155171
// fires the metrics hook. Caller MUST hold l.mu.
156172
//
157-
// The TimestampMs field is truncated from the 64-bit UnixMilli of
158-
// the sender, so the RTT calculation is done in the same truncated
159-
// space: `now.UnixMilli() & 0xFFFFFFFF` minus the received timestamp,
160-
// interpreted as a signed 32-bit delta so a clock that wraps past
161-
// 2^32 ms (every ~49 days) stays correct across the boundary.
173+
// The ring lookup is the authoritative RTT source: Machine tracks
174+
// the monotonic sentAt for every outstanding echo, so the delta is
175+
// immune to wall-clock jumps and system-suspend artifacts. When the
176+
// ring has already evicted the matching entry (burst loss, ring
177+
// overflow) the calculation falls back to the self-carried
178+
// TimestampMs in the ZEEC envelope, which is truncated from the
179+
// 64-bit UnixMilli and treated as a signed 32-bit delta so a clock
180+
// crossing the 2^32 ms boundary (~49 days) stays correct.
162181
func (l *Loop) recordEchoRTTLocked(entry *sessionEntry, e packet.Echo, now time.Time, mode string) {
163-
nowMs := uint32(now.UnixMilli())
164-
delta := int32(nowMs - e.TimestampMs)
165-
delta = max(delta, 0)
166-
rtt := time.Duration(delta) * time.Millisecond
182+
rtt, ok := entry.machine.MatchEchoRx(e.Sequence, now)
183+
if !ok {
184+
nowMs := uint32(now.UnixMilli())
185+
delta := int32(nowMs - e.TimestampMs)
186+
delta = max(delta, 0)
187+
rtt = time.Duration(delta) * time.Millisecond
188+
}
167189
entry.machine.RecordEchoRTT(rtt)
168190
if hook := l.metricsHook.Load(); hook != nil {
169191
(*hook).OnEchoRx(mode)

‎internal/plugins/bfd/engine/echo_test.go‎

Lines changed: 120 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -105,10 +105,16 @@ func TestEchoRoundTrip(t *testing.T) {
105105
if err != nil {
106106
t.Fatalf("loopB.EnsureSession: %v", err)
107107
}
108+
// Intentionally do NOT defer Unsubscribe on subA/subB: Loop.Stop
109+
// closes every subscriber channel before this test returns, and
110+
// the engine's notify path copies the subscriber slice out from
111+
// under subsMu before sending, so calling Unsubscribe AFTER the
112+
// defer chain has already started races the final echo-driven
113+
// state transition against Machine.EchoFail's notify.
114+
// TODO(spec-bfd-engine-unsubscribe-lock): tighten Unsubscribe so
115+
// it is safe to call while the express loop is still running.
108116
subA := hA.Subscribe()
109117
subB := hB.Subscribe()
110-
defer hA.Unsubscribe(subA)
111-
defer hB.Unsubscribe(subB)
112118

113119
// The three-way handshake is required for both machines to reach Up
114120
// so echoTickLocked stops clearing the echo schedule.
@@ -151,3 +157,115 @@ func TestEchoRoundTrip(t *testing.T) {
151157
t.Fatalf("no RTT samples recorded on hookA")
152158
}
153159
}
160+
161+
// dropEchoTransport is a Transport that accepts Sends silently and
162+
// never delivers any inbound packets. Used by
163+
// TestEchoDetectionSwitchover to simulate a peer that is dropping
164+
// every echo packet so the engine's echo detection path fires.
165+
//
166+
// Start / Stop are idempotent. The RX channel is never closed mid-test
167+
// so the express loop keeps selecting on a live channel; Stop closes
168+
// it so Loop.run's range exits cleanly.
169+
type dropEchoTransport struct {
170+
rx chan transport.Inbound
171+
}
172+
173+
func newDropEchoTransport() *dropEchoTransport {
174+
return &dropEchoTransport{rx: make(chan transport.Inbound)}
175+
}
176+
177+
func (*dropEchoTransport) Start() error { return nil }
178+
func (*dropEchoTransport) Send(_ transport.Outbound) error { return nil }
179+
func (t *dropEchoTransport) RX() <-chan transport.Inbound { return t.rx }
180+
func (t *dropEchoTransport) Stop() error {
181+
select {
182+
case <-t.rx:
183+
default:
184+
close(t.rx)
185+
}
186+
return nil
187+
}
188+
189+
// VALIDATES: spec-bfd-6c Phase B.2 -- when echo is negotiated and
190+
// Up but no reflected echoes return within DetectMult * EchoInterval,
191+
// the engine declares the session Down with DiagEchoFailed.
192+
// PREVENTS: regression where EchoFail never fires, fires with the
193+
// wrong diagnostic, or fails to clear the outstanding ring on
194+
// teardown.
195+
func TestEchoDetectionSwitchover(t *testing.T) {
196+
addrAA := netip.MustParseAddr(addrA)
197+
addrBB := netip.MustParseAddr(addrB)
198+
199+
controlA, controlB := transport.Pair(api.SingleHop, addrAA, addrBB)
200+
loopA := NewLoopWithEcho(controlA, newDropEchoTransport(), clock.RealClock{})
201+
loopB := NewLoopWithEcho(controlB, newDropEchoTransport(), clock.RealClock{})
202+
203+
if err := loopA.Start(); err != nil {
204+
t.Fatalf("loopA.Start: %v", err)
205+
}
206+
defer func() {
207+
if err := loopA.Stop(); err != nil {
208+
t.Errorf("loopA.Stop: %v", err)
209+
}
210+
}()
211+
if err := loopB.Start(); err != nil {
212+
t.Fatalf("loopB.Start: %v", err)
213+
}
214+
defer func() {
215+
if err := loopB.Stop(); err != nil {
216+
t.Errorf("loopB.Stop: %v", err)
217+
}
218+
}()
219+
220+
hA, err := loopA.EnsureSession(echoReqFor(addrB, addrA))
221+
if err != nil {
222+
t.Fatalf("loopA.EnsureSession: %v", err)
223+
}
224+
if _, err := loopB.EnsureSession(echoReqFor(addrA, addrB)); err != nil {
225+
t.Fatalf("loopB.EnsureSession: %v", err)
226+
}
227+
228+
// Wait for the handshake to reach Up. Subsequent transitions
229+
// (Up -> Down with DiagEchoFailed) are observed via subA.
230+
subA := hA.Subscribe()
231+
deadline := time.Now().Add(6 * time.Second)
232+
upA := false
233+
for !upA {
234+
if time.Now().After(deadline) {
235+
t.Fatalf("handshake not Up before deadline")
236+
}
237+
select {
238+
case ev := <-subA:
239+
if ev.State == packet.StateUp {
240+
upA = true
241+
}
242+
case <-time.After(50 * time.Millisecond):
243+
}
244+
}
245+
246+
// After Up the drop transport swallows every outbound echo so
247+
// the per-session outstanding ring fills up without any matching
248+
// reflections. The engine tick then observes
249+
// EchoDetectionExpired and calls EchoFail which publishes a
250+
// Down event with DiagEchoFailed.
251+
//
252+
// Detection time = DetectMult (3) * EchoInterval (10ms) = 30ms;
253+
// the poll interval is 5ms so the first expiry observation
254+
// happens ~35ms after the first echo TX. A 2-second deadline is
255+
// generously above that floor to absorb scheduler jitter under
256+
// load.
257+
deadline = time.Now().Add(2 * time.Second)
258+
gotEchoFailed := false
259+
for !gotEchoFailed {
260+
if time.Now().After(deadline) {
261+
t.Fatalf("no echo-failed transition within 2s")
262+
}
263+
select {
264+
case ev := <-subA:
265+
if ev.State == packet.StateDown && ev.Diag == packet.DiagEchoFailed {
266+
gotEchoFailed = true
267+
}
268+
case <-time.After(20 * time.Millisecond):
269+
}
270+
}
271+
}

‎internal/plugins/bfd/session/session.go‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -145,6 +145,36 @@ type Machine struct {
145145
// Zero until the first reflected echo is matched. Exposed via
146146
// LastEchoRTT for snapshot consumers.
147147
lastEchoRTT time.Duration
148+
149+
// echoOutstanding is a fixed-size ring of echo TX entries that
150+
// have not yet been matched by a returning reflection. The ring
151+
// is the state that turns echo from a passive RTT probe into an
152+
// active liveness channel: the engine declares the session Down
153+
// with DiagEchoFailed when the oldest unreturned entry has been
154+
// waiting longer than DetectMult * EchoInterval. Stage 6b Phase B.
155+
//
156+
// Slot semantics: sentAt.IsZero() marks an empty slot. The ring
157+
// is sized at echoOutstandingCap which is generous for the
158+
// default DetectMult=3 (the detection threshold is six outstanding
159+
// entries in the worst case). A full ring overwrites the oldest
160+
// slot and the overwritten entry is treated as a miss.
161+
echoOutstanding [echoOutstandingCap]echoEntry
162+
}
163+
164+
// echoOutstandingCap is the fixed ring size for per-session echo
165+
// TX tracking. Chosen to comfortably exceed 2 * DetectMult for the
166+
// common DetectMult=3 without making the Machine struct fat; the
167+
// current cap of 16 uses 16 * 24 = 384 bytes per session on a
168+
// 64-bit build.
169+
const echoOutstandingCap = 16
170+
171+
// echoEntry is one outstanding echo TX slot. sentAt is the monotonic
172+
// time the engine handed the packet to the transport; sequence is
173+
// the ZEEC Sequence field for the matching RX lookup. A zero sentAt
174+
// marks the slot as empty (swept by MatchEchoRx or ClearEchoSchedule).
175+
type echoEntry struct {
176+
sequence uint32
177+
sentAt time.Time
148178
}
149179

150180
// Role is the BFD role (Active or Passive) per RFC 5883 Section 4.3.

0 commit comments

Comments
 (0)