|
|
@@ -0,0 +1,326 @@
|
|
|
+--[[
|
|
|
+ state.lua — Watchdog state machine with hysteresis and flap detection.
|
|
|
+ Pure Lua 5.1, no dependencies.
|
|
|
+
|
|
|
+ Port of balancer-lite's internal/watchdog + internal/daemon flap logic.
|
|
|
+ Target: OpenWrt 22.03 (Lua 5.1 compat).
|
|
|
+
|
|
|
+ States:
|
|
|
+ INIT — startup, no data yet
|
|
|
+ WAN_A_PRIMARY — wan-a is active and healthy
|
|
|
+ WAN_B_PRIMARY — wan-b is active and healthy
|
|
|
+ SWITCHING_TO_A — wan-a just became healthy; switching back
|
|
|
+ SWITCHING_TO_B — wan-b just became healthy; switching back
|
|
|
+ BOTH_DOWN — both wans failed their thresholds
|
|
|
+ DEGRADED — flap detected (too many switches in window)
|
|
|
+
|
|
|
+ The machine works per-drive-cycle:
|
|
|
+ 1. Feed in probe results: feed(wan, {ok, rtt_ms, loss_pct, ...})
|
|
|
+ 2. Advance: advance() — updates windows, runs hysteresis, emits events
|
|
|
+ 3. Query: current_state(), active_wan(), last_event()
|
|
|
+
|
|
|
+ Events emitted (for the store + outbox):
|
|
|
+ { type="state.init", ts, state, active_wan, ... }
|
|
|
+ { type="failover.switch", ts, from_state, to_state, active_wan, reason }
|
|
|
+ { type="state.recovered", ts, state, active_wan }
|
|
|
+ { type="state.both_down", ts, state }
|
|
|
+ { type="state.switch_failed",ts, state, from_state, to_state, reason }
|
|
|
+ { type="flap.alert", ts, state, flap_count, window_s }
|
|
|
+]]
|
|
|
+
|
|
|
+local state = {
|
|
|
+ INIT = "INIT",
|
|
|
+ WAN_A_PRIMARY = "WAN_A_PRIMARY",
|
|
|
+ WAN_B_PRIMARY = "WAN_B_PRIMARY",
|
|
|
+ SWITCHING_TO_A = "SWITCHING_TO_A",
|
|
|
+ SWITCHING_TO_B = "SWITCHING_TO_B",
|
|
|
+ BOTH_DOWN = "BOTH_DOWN",
|
|
|
+ DEGRADED = "DEGRADED",
|
|
|
+}
|
|
|
+
|
|
|
+local function new(cfg)
|
|
|
+ cfg = cfg or {}
|
|
|
+ local T = {
|
|
|
+ -- Config (passed from UCI)
|
|
|
+ wan_a_id = cfg.wan_a_id or "wan-a",
|
|
|
+ wan_b_id = cfg.wan_b_id or "wan-b",
|
|
|
+ down_thr = cfg.down_threshold or 5, -- consecutive fails → DOWN
|
|
|
+ up_thr = cfg.up_threshold or 10, -- consecutive ok → UP
|
|
|
+ flap_thr = cfg.flap_threshold or 3, -- switches in window → DEGRADED
|
|
|
+ flap_window_s = cfg.flap_window_s or 300, -- flap detection window (5 min)
|
|
|
+ flap_recovery = cfg.flap_recovery or 300, -- must stay calm this long to exit DEGRADED
|
|
|
+
|
|
|
+ -- Per-WAN sliding window (ring buffer of last window_size results)
|
|
|
+ window_size = cfg.window_size or 10,
|
|
|
+ wan_a_window = {}, -- {ok=true/false, rtt_ms, ts}
|
|
|
+ wan_b_window = {},
|
|
|
+
|
|
|
+ -- Streak counters (consecutive failures / ok per WAN)
|
|
|
+ wan_a_streak = 0, -- positive=ok streak, negative=fail streak
|
|
|
+ wan_b_streak = 0,
|
|
|
+
|
|
|
+ -- Flap detection
|
|
|
+ switch_times = {}, -- timestamps of recent switches
|
|
|
+
|
|
|
+ -- Current state
|
|
|
+ cur_state = state.INIT,
|
|
|
+ active_wan = nil,
|
|
|
+ last_event = nil, -- last emitted event
|
|
|
+
|
|
|
+ -- Tracking
|
|
|
+ degraded_since = nil, -- timestamp when DEGRADED was entered
|
|
|
+ last_switch_ts = 0,
|
|
|
+ }
|
|
|
+ return T
|
|
|
+end
|
|
|
+
|
|
|
+-- Mark a probe result for a WAN
|
|
|
+local function feed(T, wan_id, result)
|
|
|
+ local win = (wan_id == T.wan_a_id) and T.wan_a_window or T.wan_b_window
|
|
|
+ local streak_key = (wan_id == T.wan_a_id) and "wan_a_streak" or "wan_b_streak"
|
|
|
+
|
|
|
+ table.insert(win, { ok = result.ok, rtt_ms = result.rtt_ms, ts = result.ts or os.time() })
|
|
|
+ if #win > T.window_size then table.remove(win, 1) end
|
|
|
+
|
|
|
+ -- Update streak: positive = ok streak, negative = fail streak
|
|
|
+ if result.ok then
|
|
|
+ T[streak_key] = math.max(1, T[streak_key] + 1)
|
|
|
+ else
|
|
|
+ T[streak_key] = math.min(-1, T[streak_key] - 1)
|
|
|
+ end
|
|
|
+end
|
|
|
+
|
|
|
+-- How many consecutive ok probes does a WAN have right now?
|
|
|
+local function ok_streak(T, wan_id)
|
|
|
+ local win = (wan_id == T.wan_a_id) and T.wan_a_window or T.wan_b_window
|
|
|
+ local streak = 0
|
|
|
+ for i = #win, 1, -1 do
|
|
|
+ if win[i].ok then streak = streak + 1
|
|
|
+ else break end
|
|
|
+ end
|
|
|
+ return streak
|
|
|
+end
|
|
|
+
|
|
|
+-- How many consecutive fail probes does a WAN have right now?
|
|
|
+local function fail_streak(T, wan_id)
|
|
|
+ local win = (wan_id == T.wan_a_id) and T.wan_a_window or T.wan_b_window
|
|
|
+ local streak = 0
|
|
|
+ for i = #win, 1, -1 do
|
|
|
+ if not win[i].ok then streak = streak + 1
|
|
|
+ else break end
|
|
|
+ end
|
|
|
+ return streak
|
|
|
+end
|
|
|
+
|
|
|
+local function is_up(T, wan_id)
|
|
|
+ return ok_streak(T, wan_id) >= T.up_thr
|
|
|
+end
|
|
|
+
|
|
|
+local function is_down(T, wan_id)
|
|
|
+ return fail_streak(T, wan_id) >= T.down_thr
|
|
|
+end
|
|
|
+
|
|
|
+-- Emit an event (stores in T.last_event; caller should copy to outbox)
|
|
|
+local function emit(T, ev)
|
|
|
+ ev.ts = ev.ts or os.time()
|
|
|
+ T.last_event = ev
|
|
|
+ return ev
|
|
|
+end
|
|
|
+
|
|
|
+-- Record a switch timestamp for flap detection
|
|
|
+local function record_switch(T)
|
|
|
+ local now = os.time()
|
|
|
+ table.insert(T.switch_times, now)
|
|
|
+ T.last_switch_ts = now
|
|
|
+ -- Prune old entries outside flap_window
|
|
|
+ local cutoff = now - T.flap_window_s
|
|
|
+ while T.switch_times[1] and T.switch_times[1] < cutoff do
|
|
|
+ table.remove(T.switch_times, 1)
|
|
|
+ end
|
|
|
+end
|
|
|
+
|
|
|
+-- Advance the state machine one drive cycle
|
|
|
+-- Returns: { state, active_wan, event } or nil if no change
|
|
|
+local function advance(T)
|
|
|
+ local prev_state = T.cur_state
|
|
|
+ local prev_active = T.active_wan
|
|
|
+ local now = os.time()
|
|
|
+
|
|
|
+ -- Flap recovery: if DEGRADED and calm long enough, exit DEGRADED
|
|
|
+ if T.cur_state == state.DEGRADED and T.degraded_since then
|
|
|
+ if (now - T.degraded_since) >= T.flap_recovery then
|
|
|
+ T.cur_state = state.INIT
|
|
|
+ T.degraded_since = nil
|
|
|
+ return emit(T, { type = "state.recovered", ts = now,
|
|
|
+ state = state.INIT, active_wan = nil })
|
|
|
+ end
|
|
|
+ end
|
|
|
+
|
|
|
+ -- If INIT, try to promote to whichever WAN is up
|
|
|
+ if T.cur_state == state.INIT then
|
|
|
+ local a_up = is_up(T, T.wan_a_id)
|
|
|
+ local b_up = is_up(T, T.wan_b_id)
|
|
|
+ if a_up and b_up then
|
|
|
+ -- Both up — prefer higher preference (A by default)
|
|
|
+ T.cur_state = state.WAN_A_PRIMARY; T.active_wan = T.wan_a_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "state.init", ts = now, state = state.WAN_A_PRIMARY,
|
|
|
+ active_wan = T.wan_a_id })
|
|
|
+ elseif a_up then
|
|
|
+ T.cur_state = state.WAN_A_PRIMARY; T.active_wan = T.wan_a_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "state.init", ts = now, state = state.WAN_A_PRIMARY,
|
|
|
+ active_wan = T.wan_a_id })
|
|
|
+ elseif b_up then
|
|
|
+ T.cur_state = state.WAN_B_PRIMARY; T.active_wan = T.wan_b_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "state.init", ts = now, state = state.WAN_B_PRIMARY,
|
|
|
+ active_wan = T.wan_b_id })
|
|
|
+ end
|
|
|
+ -- Stay INIT (no WAN up yet — keep probing)
|
|
|
+ return nil
|
|
|
+ end
|
|
|
+
|
|
|
+ -- Flap detection check (not in INIT, SWITCHING, or BOTH_DOWN)
|
|
|
+ local is_transitional = (T.cur_state == state.SWITCHING_TO_A or
|
|
|
+ T.cur_state == state.SWITCHING_TO_B)
|
|
|
+ if not is_transitional and T.cur_state ~= state.BOTH_DOWN then
|
|
|
+ local nswitches = #T.switch_times
|
|
|
+ if nswitches >= T.flap_thr then
|
|
|
+ -- Enter DEGRADED
|
|
|
+ T.cur_state = state.DEGRADED
|
|
|
+ T.degraded_since = now
|
|
|
+ return emit(T, { type = "flap.alert", ts = now,
|
|
|
+ state = state.DEGRADED, flap_count = nswitches,
|
|
|
+ window_s = T.flap_window_s })
|
|
|
+ end
|
|
|
+ end
|
|
|
+
|
|
|
+ -- BOTH_DOWN: wait for either WAN to recover
|
|
|
+ if T.cur_state == state.BOTH_DOWN then
|
|
|
+ local a_up = is_up(T, T.wan_a_id)
|
|
|
+ local b_up = is_up(T, T.wan_b_id)
|
|
|
+ if a_up and b_up then
|
|
|
+ T.cur_state = state.WAN_A_PRIMARY; T.active_wan = T.wan_a_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "state.recovered", ts = now,
|
|
|
+ from_state = state.BOTH_DOWN,
|
|
|
+ state = state.WAN_A_PRIMARY, active_wan = T.wan_a_id })
|
|
|
+ elseif a_up then
|
|
|
+ T.cur_state = state.WAN_A_PRIMARY; T.active_wan = T.wan_a_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "state.recovered", ts = now,
|
|
|
+ from_state = state.BOTH_DOWN,
|
|
|
+ state = state.WAN_A_PRIMARY, active_wan = T.wan_a_id })
|
|
|
+ elseif b_up then
|
|
|
+ T.cur_state = state.WAN_B_PRIMARY; T.active_wan = T.wan_b_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "state.recovered", ts = now,
|
|
|
+ from_state = state.BOTH_DOWN,
|
|
|
+ state = state.WAN_B_PRIMARY, active_wan = T.wan_b_id })
|
|
|
+ end
|
|
|
+ return nil
|
|
|
+ end
|
|
|
+
|
|
|
+ -- Normal states (WAN_A_PRIMARY, WAN_B_PRIMARY, DEGRADED)
|
|
|
+ local primary = T.active_wan
|
|
|
+ local standby = (primary == T.wan_a_id) and T.wan_b_id or T.wan_a_id
|
|
|
+
|
|
|
+ local primary_down = is_down(T, primary)
|
|
|
+ local standby_up = is_up(T, standby)
|
|
|
+ local standby_down = is_down(T, standby)
|
|
|
+
|
|
|
+ -- Case 1: Primary failed, standby is up → switch
|
|
|
+ if primary_down and standby_up then
|
|
|
+ local to_state = (primary == T.wan_a_id) and state.SWITCHING_TO_B or state.SWITCHING_TO_A
|
|
|
+ local new_state = (standby == T.wan_a_id) and state.WAN_A_PRIMARY or state.WAN_B_PRIMARY
|
|
|
+ T.cur_state = to_state
|
|
|
+ T.active_wan = standby
|
|
|
+ record_switch(T)
|
|
|
+ -- After a brief settling period the state will be finalised by the
|
|
|
+ -- caller (in the real daemon this is done by rtctl after the switch
|
|
|
+ -- completes). Here we emit the transition event.
|
|
|
+ return emit(T, { type = "failover.switch", ts = now,
|
|
|
+ from_state = prev_state, to_state = new_state,
|
|
|
+ active_wan = standby,
|
|
|
+ reason = primary .. ": " .. tostring(T.down_thr) ..
|
|
|
+ "/" .. tostring(T.down_thr) .. " probes failed" })
|
|
|
+
|
|
|
+ -- Case 2: Primary failed, standby also failed → BOTH_DOWN
|
|
|
+ elseif primary_down and standby_down then
|
|
|
+ T.cur_state = state.BOTH_DOWN
|
|
|
+ T.active_wan = nil
|
|
|
+ return emit(T, { type = "state.both_down", ts = now, state = state.BOTH_DOWN })
|
|
|
+
|
|
|
+ -- Case 3: In SWITCHING_TO_A/B — standby confirmed up → finalise
|
|
|
+ elseif T.cur_state == state.SWITCHING_TO_A and is_up(T, T.wan_a_id) then
|
|
|
+ T.cur_state = state.WAN_A_PRIMARY
|
|
|
+ T.active_wan = T.wan_a_id
|
|
|
+ return nil -- no new event on finalisation
|
|
|
+
|
|
|
+ elseif T.cur_state == state.SWITCHING_TO_B and is_up(T, T.wan_b_id) then
|
|
|
+ T.cur_state = state.WAN_B_PRIMARY
|
|
|
+ T.active_wan = T.wan_b_id
|
|
|
+ return nil
|
|
|
+
|
|
|
+ -- Case 4: Primary recovered while in SWITCHING — abort and return
|
|
|
+ -- (handled above via the SWITCHING_TO_* finalisation)
|
|
|
+
|
|
|
+ -- Case 5: Primary recovered while in DEGRADED — stay degraded until flap clears
|
|
|
+ elseif not primary_down and T.cur_state == state.DEGRADED then
|
|
|
+ -- stay DEGRADED, but record that primary is now healthy
|
|
|
+ return nil
|
|
|
+
|
|
|
+ -- Case 6: Preferred WAN (A) recovered while on B — initiate switch-back
|
|
|
+ -- Only when not in DEGRADED or SWITCHING state
|
|
|
+ elseif T.cur_state == state.WAN_B_PRIMARY and
|
|
|
+ not is_transitional and
|
|
|
+ is_up(T, T.wan_a_id) and
|
|
|
+ primary ~= T.wan_a_id then
|
|
|
+ -- wan-a (preferred) recovered → switch back
|
|
|
+ T.cur_state = state.SWITCHING_TO_A
|
|
|
+ T.active_wan = T.wan_a_id
|
|
|
+ record_switch(T)
|
|
|
+ return emit(T, { type = "failover.switch", ts = now,
|
|
|
+ from_state = state.WAN_B_PRIMARY,
|
|
|
+ to_state = state.WAN_A_PRIMARY,
|
|
|
+ active_wan = T.wan_a_id,
|
|
|
+ reason = T.wan_a_id .. " recovered" })
|
|
|
+
|
|
|
+ -- Case 7: wan-b recovered while on A and A is still up — stay on A
|
|
|
+ -- (deliberate: prefer primary, only switch on primary failure)
|
|
|
+ end
|
|
|
+
|
|
|
+ return nil
|
|
|
+end
|
|
|
+
|
|
|
+-- Return verdict info for a specific WAN (for metrics / store)
|
|
|
+local function verdict(T, wan_id)
|
|
|
+ local win = (wan_id == T.wan_a_id) and T.wan_a_window or T.wan_b_window
|
|
|
+ local total = #win
|
|
|
+ if total == 0 then return { verdict = "unknown", loss_pct = 0, rtt_ms = nil } end
|
|
|
+ local failed = 0; local rtt_sum = 0; local rtt_n = 0
|
|
|
+ for i = 1, total do
|
|
|
+ if not win[i].ok then failed = failed + 1
|
|
|
+ else
|
|
|
+ if win[i].rtt_ms then rtt_sum = rtt_sum + win[i].rtt_ms; rtt_n = rtt_n + 1 end
|
|
|
+ end
|
|
|
+ end
|
|
|
+ local loss_pct = (failed / total) * 100
|
|
|
+ local rtt_avg = rtt_n > 0 and (rtt_sum / rtt_n) or nil
|
|
|
+ local verdict
|
|
|
+ if fail_streak(T, wan_id) >= T.down_thr then verdict = "down"
|
|
|
+ elseif ok_streak(T, wan_id) >= T.up_thr then verdict = "healthy"
|
|
|
+ else verdict = "degraded" end
|
|
|
+ return { verdict = verdict, loss_pct = loss_pct, rtt_ms = rtt_avg,
|
|
|
+ failed = failed, total = total }
|
|
|
+end
|
|
|
+
|
|
|
+-- Module
|
|
|
+return {
|
|
|
+ new = new,
|
|
|
+ feed = feed,
|
|
|
+ advance = advance,
|
|
|
+ verdict = verdict,
|
|
|
+ state = state,
|
|
|
+}
|