From 15fa6bd229b008e845d193caf1a8969a150d726c Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Wed, 8 Jul 2026 12:34:53 -0400 Subject: [PATCH 1/9] Add TCP adaptor connection gathering for idle-connection sweeper Cross-references the router's TCP adaptor connections (skmanage QUERY) with kernel socket state (ss -tin) to get an idle-time signal more accurate than the router's own lastDlvSeconds, which only reflects delivery-object creation and not ongoing byte traffic on a long-lived stream. Includes gatherdump, a standalone tool for inspecting a Snapshot against a live router. --- internal/cmd/skupper/debug/sweeper/gather.go | 142 ++++++++++++++++++ .../skupper/debug/sweeper/gatherdump/main.go | 41 +++++ 2 files changed, 183 insertions(+) create mode 100644 internal/cmd/skupper/debug/sweeper/gather.go create mode 100644 internal/cmd/skupper/debug/sweeper/gatherdump/main.go diff --git a/internal/cmd/skupper/debug/sweeper/gather.go b/internal/cmd/skupper/debug/sweeper/gather.go new file mode 100644 index 000000000..fd1a520ba --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/gather.go @@ -0,0 +1,142 @@ +package sweeper + +import ( + "encoding/json" + "fmt" + "os/exec" + "strconv" + "strings" + "time" +) + +const ( + ConnType = "io.skupper.router.connection" + tcpContainer = "TcpAdaptor" + egressDispatch = "egress-dispatch" +) + +// connInfo is the router's view of a single connection, as returned by +// `skmanage QUERY --type=io.skupper.router.connection`. +type connInfo struct { + Identity string `json:"identity"` + Container string `json:"container"` + Host string `json:"host"` + Dir string `json:"dir"` + UptimeSeconds *int `json:"uptimeSeconds"` + LastDlvSeconds *int `json:"lastDlvSeconds"` +} + +// socketInfo is the kernel's view of a TCP socket, as reported by `ss -tin`. +// LastRcvMs/LastSndMs come straight from TCP_INFO and should show actual activity vs lastDlvSeconds +type socketInfo struct { + LastRcvMs int + LastSndMs int +} + +// Snapshot bundles the router's connection list and the kernel's socket +// state at one point in time. +type Snapshot struct { + Now time.Time + TCPConns []connInfo + // Sockets is keyed by peer "host:port" so it can be matched against an + // inbound connInfo.Host. Only sockets for 'in' (client-facing) connections + // can be matched reliably. + Sockets map[string]socketInfo +} + +// Gather queries the router for its TCP adaptor connections and cross +// references them with kernel socket state. It performs no filtering beyond +// discarding non-TCP-adaptor connections, and makes no decisions about which +// connections are healthy. +func Gather(skmanageBin, url string) (Snapshot, error) { + raw, err := runSkmanage(skmanageBin, url, "QUERY", "--type="+ConnType) + if err != nil { + return Snapshot{}, fmt.Errorf("could not query router at %s: %w", url, err) + } + + var allConns []connInfo + if err := json.Unmarshal(raw, &allConns); err != nil { + return Snapshot{}, fmt.Errorf("failed to parse connection list: %w", err) + } + + var tcpConns []connInfo + for _, c := range allConns { + if isTCPAdaptorConn(c) { + tcpConns = append(tcpConns, c) + } + } + + return Snapshot{ + Now: time.Now(), + TCPConns: tcpConns, + Sockets: gatherSockets(), + }, nil +} + +func isTCPAdaptorConn(c connInfo) bool { + return c.Container == tcpContainer && c.Host != egressDispatch +} + +//runs ss -tin and builds a peer-address → {lastrcv, lastsnd} map by pairing each socket's +//header line with its following detail line. +func gatherSockets() map[string]socketInfo { + sockets := map[string]socketInfo{} + out, err := exec.Command("ss", "-tin").Output() + if err != nil { + return sockets + } + + var pendingPeer string + for _, line := range strings.Split(string(out), "\n") { + if line == "" { + continue + } + if line[0] != ' ' && line[0] != '\t' { + pendingPeer = "" + fields := strings.Fields(line) + if len(fields) < 5 || fields[0] == "State" { + continue + } + pendingPeer = fields[4] + continue + } + if pendingPeer == "" { + continue + } + sockets[pendingPeer] = socketInfo{ + LastRcvMs: extractMsField(line, "lastrcv:"), + LastSndMs: extractMsField(line, "lastsnd:"), + } + pendingPeer = "" + } + return sockets +} + +// extractMsField returns the integer following "key:" in line (e.g. key +// "lastrcv:" on "... lastrcv:592 lastack:9119 ..." returns 592). Returns 0 +// if key isn't found in line. +func extractMsField(line, key string) int { + idx := strings.Index(line, key) + if idx == -1 { + return 0 + } + rest := line[idx+len(key):] + end := strings.IndexAny(rest, " \t") + if end != -1 { + rest = rest[:end] + } + val, err := strconv.Atoi(rest) + if err != nil { + return 0 + } + return val +} + +func runSkmanage(bin, url string, args ...string) ([]byte, error) { + cmdArgs := append([]string{"--bus", url}, args...) + out, err := exec.Command(bin, cmdArgs...).Output() + if err != nil { + return nil, fmt.Errorf("skmanage failed: %w", err) + } + return out, nil +} diff --git a/internal/cmd/skupper/debug/sweeper/gatherdump/main.go b/internal/cmd/skupper/debug/sweeper/gatherdump/main.go new file mode 100644 index 000000000..909f28c78 --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/gatherdump/main.go @@ -0,0 +1,41 @@ +// Throwaway debug tool: calls sweeper.Gather directly and dumps the results for debugging +package main + +import ( + "flag" + "fmt" + "os" + + "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" +) + +func main() { + url := flag.String("url", "amqp://127.0.0.1:5672", "Router management URL") + skmanageBin := flag.String("skmanage", "skmanage", "Path to skmanage binary") + flag.Parse() + + snap, err := sweeper.Gather(*skmanageBin, *url) + if err != nil { + fmt.Fprintln(os.Stderr, "gather failed:", err) + os.Exit(1) + } + + fmt.Printf("gathered at %s from %s\n\n", snap.Now.Format("15:04:05"), *url) + fmt.Printf("%-4s %-25s %-10s %-10s\n", "DIR", "HOST", "UPTIME(s)", "LASTDLV(s)") + for _, c := range snap.TCPConns { + fmt.Printf("%-4s %-25s %-10s %-10s\n", + c.Dir, c.Host, ptrStr(c.UptimeSeconds), ptrStr(c.LastDlvSeconds)) + } + + fmt.Println("\nsockets (host:port -> lastrcv/lastsnd s):") + for host, s := range snap.Sockets { + fmt.Printf(" %-25s lastrcv=%.1fs lastsnd=%.1fs\n", host, float64(s.LastRcvMs)/1000, float64(s.LastSndMs)/1000) + } +} + +func ptrStr(p *int) string { + if p == nil { + return "-" + } + return fmt.Sprintf("%d", *p) +} From a83615cd53a29c4810fbae9e3940c032aafc0beb Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Thu, 9 Jul 2026 12:10:50 -0400 Subject: [PATCH 2/9] Added go script to display TCP connections sending / receiving Info to one another --- internal/cmd/skupper/debug/sweeper/README.md | 284 ++++++++++++++++++ internal/cmd/skupper/debug/sweeper/gather.go | 84 ++++-- .../skupper/debug/sweeper/gatherdump/main.go | 15 +- 3 files changed, 349 insertions(+), 34 deletions(-) create mode 100644 internal/cmd/skupper/debug/sweeper/README.md diff --git a/internal/cmd/skupper/debug/sweeper/README.md b/internal/cmd/skupper/debug/sweeper/README.md new file mode 100644 index 000000000..3a7e82194 --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/README.md @@ -0,0 +1,284 @@ +# Connection sweeper — data gathering demo + +## Setup + +Needs: a built `skupper-router`, plus the `a.conf`, `b.conf`, and `random_traffic_sim.py` files below. + +Create/modify these three files in order to run the tests. + +--- + +`a.conf` — in `/build/` + +```conf +router { + mode: interior + id: A +} + +listener { + port: amqp + authenticatePeer: no + saslMechanisms: ANONYMOUS +} + +listener { + port: 10000 + authenticatePeer: no + saslMechanisms: ANONYMOUS + role: inter-router +} + +tcpListener { + port: 9090 + address: tcp-echo + name: tcp-echo-listener +} +``` + +--- + +`b.conf` — in `/build/` + +```conf +router { + mode: interior + id: B +} +listener { + host: 127.0.0.1 + port: 5673 + role: normal +} +connector { + host: 127.0.0.1 + port: 10000 + role: inter-router +} +tcpConnector { + address: tcp-echo + host: 127.0.0.1 + port: 9091 + name: tcp-echo-connector +} +``` + +--- + +`random_traffic_sim.py` — in `/scripts/` + +```python +#!/usr/bin/env python3 +""" +Opens random TCP connections to a router's TCP adaptor, each sending a small payload at a random interval. +""" + +import argparse +import random +import signal +import socket +import threading +import time +from datetime import datetime + + +def log(msg: str) -> None: + ts = datetime.now().strftime("%H:%M:%S") + print(f"[{ts}] {msg}", flush=True) + + +def _run_backend(port: int) -> None: + srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + srv.bind(("0.0.0.0", port)) + srv.listen(256) + while True: + try: + conn, _ = srv.accept() + threading.Thread(target=_backend_echo, args=(conn,), daemon=True).start() + except OSError: + break + + +def _backend_echo(conn: socket.socket) -> None: + try: + while True: + data = conn.recv(4096) + if not data: + return + conn.sendall(data) + except OSError: + pass + finally: + conn.close() + + +def _connection_worker(idx: int, host: str, port: int, min_interval: int, max_interval: int, + stop: threading.Event) -> None: + """Holds one TCP connection open for the whole run, sending a payload and + reading the echo back at an independently randomized interval.""" + sock = None + try: + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.connect((host, port)) + sock.sendall(f"conn-{idx} open\n".encode()) + sock.settimeout(5.0) + try: + sock.recv(4096) + except socket.timeout: + pass + + while not stop.is_set(): + interval = random.randint(min_interval, max_interval) + if stop.wait(interval): + break + try: + payload = f"conn-{idx} ping {int(time.time())}\n".encode() + sock.sendall(payload) + sock.recv(4096) + except OSError as e: + log(f"conn-{idx} error, exiting: {e}") + return + except OSError as e: + log(f"conn-{idx} failed to connect: {e}") + finally: + if sock: + sock.close() + + +def main() -> None: + parser = argparse.ArgumentParser( + prog="random_traffic_sim.py", + description="Steady TCP connections sending data at random intervals, for sweeper testing.", + formatter_class=argparse.RawDescriptionHelpFormatter, + epilog=__doc__, + ) + parser.add_argument("--host", default="127.0.0.1", help="Router host (default: 127.0.0.1)") + parser.add_argument("--port", type=int, default=9090, help="Router TCP adaptor port (default: 9090)") + parser.add_argument("--connections", type=int, default=20, help="Number of connections to open (default: 20)") + parser.add_argument("--min-interval", type=int, default=30, + help="Minimum seconds between sends on a connection (default: 30)") + parser.add_argument("--max-interval", type=int, default=500, + help="Maximum seconds between sends on a connection (default: 500)") + parser.add_argument("--backend-port", type=int, default=None, + help="Port for this script's own echo backend (default: --port + 1)") + args = parser.parse_args() + + backend_port = args.backend_port or (args.port + 1) + stop = threading.Event() + + def _on_sigint(_s, _f): + log("Interrupted — shutting down ...") + stop.set() + + signal.signal(signal.SIGINT, _on_sigint) + + threading.Thread(target=_run_backend, args=(backend_port,), daemon=True).start() + time.sleep(0.3) + log(f"Backend echoing on :{backend_port}") + + log(f"Opening {args.connections} connection(s) to {args.host}:{args.port}, " + f"each sending every {args.min_interval}-{args.max_interval}s ...") + threads = [] + for i in range(args.connections): + t = threading.Thread( + target=_connection_worker, + args=(i, args.host, args.port, args.min_interval, args.max_interval, stop), + daemon=True, + ) + t.start() + threads.append(t) + time.sleep(0.05) + + log("All connections open. Running until Ctrl+C.") + stop.wait() + + for t in threads: + t.join(timeout=2) + log("Exited.") + + +if __name__ == "__main__": + main() +``` + +--- + +Copy-paste, once per terminal (need 4 terminals): + +```bash +export SKUPPER_ROUTER=~/skupper-router # adjust if your checkout lives elsewhere +export SKUPPER_REPO=~/skupper +[ -f "$SKUPPER_ROUTER/build/config.sh" ] || echo "!! not found — fix SKUPPER_ROUTER above" +source "$SKUPPER_ROUTER/build/config.sh" +cd "$SKUPPER_REPO" +``` + +Structured as client -> Router A -> Router B -> echo backend + +## Run it + +**Terminal 1 — Router A:** +```bash +"$SKUPPER_ROUTER/build/router/skrouterd" -c "$SKUPPER_ROUTER/build/a.conf" +``` + +**Terminal 2 — Router B:** +```bash +"$SKUPPER_ROUTER/build/router/skrouterd" -c "$SKUPPER_ROUTER/build/b.conf" +``` + +**Terminal 3 — traffic** (starts its own echo backend on `:9091`): +```bash +python3 "$SKUPPER_ROUTER/scripts/random_traffic_sim.py" --connections 6 --min-interval 30 --max-interval 300 +``` + +Wait for `All connections open.`, then run the tests below in a fourth terminal. + +--- + +## Test 1 — the `in` side + +```bash +go run ./internal/cmd/skupper/debug/sweeper/gatherdump --url amqp://127.0.0.1:5672 +``` + +``` +DIR HOST LOCALSOCKET UPTIME(s) LASTDLV(s) +in 127.0.0.1:53484 127.0.0.1:9090 123 123 +in 127.0.0.1:53498 127.0.0.1:9090 123 123 +in 127.0.0.1:53506 127.0.0.1:9090 123 123 +in 127.0.0.1:53514 127.0.0.1:9090 123 123 +in 127.0.0.1:53528 127.0.0.1:9090 123 123 +in 127.0.0.1:53538 127.0.0.1:9090 123 123 +``` + +`HOST` is unique per connection. `LOCALSOCKET` is the same on every row cause it's the listener. + +The two `sockets by ...` sections below the table are the same kernel sockets indexed by each end — `in` connections are matched by their peer address, `out` connections by their local address. + + +## Test 2 — the `out` side + +Same command, Router B: + +```bash +go run ./internal/cmd/skupper/debug/sweeper/gatherdump --url amqp://127.0.0.1:5673 +``` + +``` +DIR HOST LOCALSOCKET UPTIME(s) LASTDLV(s) +out 127.0.0.1:9091 127.0.0.1:54944 133 133 +out 127.0.0.1:9091 127.0.0.1:54948 133 133 +out 127.0.0.1:9091 127.0.0.1:54956 133 133 +out 127.0.0.1:9091 127.0.0.1:54960 133 133 +out 127.0.0.1:9091 127.0.0.1:54972 133 133 +out 127.0.0.1:9091 127.0.0.1:54976 133 133 +``` + +## Clean up +```bash +pkill -f 'skrouterd -c .*a\.conf' +pkill -f 'skrouterd -c .*b\.conf' +pkill -f random_traffic_sim.py +``` diff --git a/internal/cmd/skupper/debug/sweeper/gather.go b/internal/cmd/skupper/debug/sweeper/gather.go index fd1a520ba..b88fcb085 100644 --- a/internal/cmd/skupper/debug/sweeper/gather.go +++ b/internal/cmd/skupper/debug/sweeper/gather.go @@ -18,9 +18,13 @@ const ( // connInfo is the router's view of a single connection, as returned by // `skmanage QUERY --type=io.skupper.router.connection`. type connInfo struct { - Identity string `json:"identity"` - Container string `json:"container"` - Host string `json:"host"` + Identity string `json:"identity"` + Container string `json:"container"` + Host string `json:"host"` + // LocalSocket is the router's own "host:port" for this connection's + // socket. Unlike Host, it is unique per connection even for 'out' + // connections, which all share the backend's address as their Host. + LocalSocket string `json:"localSocket"` Dir string `json:"dir"` UptimeSeconds *int `json:"uptimeSeconds"` LastDlvSeconds *int `json:"lastDlvSeconds"` @@ -38,18 +42,28 @@ type socketInfo struct { type Snapshot struct { Now time.Time TCPConns []connInfo - // Sockets is keyed by peer "host:port" so it can be matched against an - // inbound connInfo.Host. Only sockets for 'in' (client-facing) connections - // can be matched reliably. + // Sockets is keyed by peer "host:port", matching an 'in' connection's + // Host (each client has a unique peer address). Sockets map[string]socketInfo + // SocketsByLocal is keyed by the socket's own local "host:port", + // matching an 'out' connection's LocalSocket ('out' peers all share the + // backend's address, so only the local side is unique). + SocketsByLocal map[string]socketInfo +} + +// Execer runs a command (argv) and returns its stdout. Both skmanage and the +// socket query go through it, so they always observe the same host — and so +// the same network namespace, which is what makes their results joinable. +type Execer func(argv []string) ([]byte, error) + +func LocalExec(argv []string) ([]byte, error) { + return exec.Command(argv[0], argv[1:]...).Output() } // Gather queries the router for its TCP adaptor connections and cross -// references them with kernel socket state. It performs no filtering beyond -// discarding non-TCP-adaptor connections, and makes no decisions about which -// connections are healthy. -func Gather(skmanageBin, url string) (Snapshot, error) { - raw, err := runSkmanage(skmanageBin, url, "QUERY", "--type="+ConnType) +// references them with kernel socket state. Discards non-TCP-adaptor connections +func Gather(execFn Execer, skmanageBin, url string) (Snapshot, error) { + raw, err := runSkmanage(execFn, skmanageBin, url, "QUERY", "--type="+ConnType) if err != nil { return Snapshot{}, fmt.Errorf("could not query router at %s: %w", url, err) } @@ -66,10 +80,12 @@ func Gather(skmanageBin, url string) (Snapshot, error) { } } + byPeer, byLocal := gatherSockets(execFn) return Snapshot{ - Now: time.Now(), - TCPConns: tcpConns, - Sockets: gatherSockets(), + Now: time.Now(), + TCPConns: tcpConns, + Sockets: byPeer, + SocketsByLocal: byLocal, }, nil } @@ -77,39 +93,49 @@ func isTCPAdaptorConn(c connInfo) bool { return c.Container == tcpContainer && c.Host != egressDispatch } -//runs ss -tin and builds a peer-address → {lastrcv, lastsnd} map by pairing each socket's -//header line with its following detail line. -func gatherSockets() map[string]socketInfo { - sockets := map[string]socketInfo{} - out, err := exec.Command("ss", "-tin").Output() +// gatherSockets reads kernel socket state via `ss -tin`. A host without ss (to be patched later) +// yields no sockets, which leaves every connection unmatched and untouched. +func gatherSockets(execFn Execer) (byPeer, byLocal map[string]socketInfo) { + out, err := execFn([]string{"ss", "-tin"}) if err != nil { - return sockets + return map[string]socketInfo{}, map[string]socketInfo{} } + return socketsFromSS(out) +} + +// socketsFromSS builds two {lastrcv, lastsnd} maps — one keyed by peer +// address, one by local address — by pairing each socket's header line in +// `ss -tin` output with its following detail line. +func socketsFromSS(out []byte) (byPeer, byLocal map[string]socketInfo) { + byPeer = map[string]socketInfo{} + byLocal = map[string]socketInfo{} - var pendingPeer string + var pendingLocal, pendingPeer string for _, line := range strings.Split(string(out), "\n") { if line == "" { continue } if line[0] != ' ' && line[0] != '\t' { - pendingPeer = "" + pendingLocal, pendingPeer = "", "" fields := strings.Fields(line) if len(fields) < 5 || fields[0] == "State" { continue } - pendingPeer = fields[4] + pendingLocal, pendingPeer = fields[3], fields[4] continue } if pendingPeer == "" { continue } - sockets[pendingPeer] = socketInfo{ + sock := socketInfo{ LastRcvMs: extractMsField(line, "lastrcv:"), LastSndMs: extractMsField(line, "lastsnd:"), } - pendingPeer = "" + byPeer[pendingPeer] = sock + byLocal[pendingLocal] = sock + pendingLocal, pendingPeer = "", "" } - return sockets + return byPeer, byLocal } // extractMsField returns the integer following "key:" in line (e.g. key @@ -132,9 +158,9 @@ func extractMsField(line, key string) int { return val } -func runSkmanage(bin, url string, args ...string) ([]byte, error) { - cmdArgs := append([]string{"--bus", url}, args...) - out, err := exec.Command(bin, cmdArgs...).Output() +func runSkmanage(execFn Execer, bin, url string, args ...string) ([]byte, error) { + argv := append([]string{bin, "--bus", url}, args...) + out, err := execFn(argv) if err != nil { return nil, fmt.Errorf("skmanage failed: %w", err) } diff --git a/internal/cmd/skupper/debug/sweeper/gatherdump/main.go b/internal/cmd/skupper/debug/sweeper/gatherdump/main.go index 909f28c78..a3115c554 100644 --- a/internal/cmd/skupper/debug/sweeper/gatherdump/main.go +++ b/internal/cmd/skupper/debug/sweeper/gatherdump/main.go @@ -14,23 +14,28 @@ func main() { skmanageBin := flag.String("skmanage", "skmanage", "Path to skmanage binary") flag.Parse() - snap, err := sweeper.Gather(*skmanageBin, *url) + snap, err := sweeper.Gather(sweeper.LocalExec, *skmanageBin, *url) if err != nil { fmt.Fprintln(os.Stderr, "gather failed:", err) os.Exit(1) } fmt.Printf("gathered at %s from %s\n\n", snap.Now.Format("15:04:05"), *url) - fmt.Printf("%-4s %-25s %-10s %-10s\n", "DIR", "HOST", "UPTIME(s)", "LASTDLV(s)") + fmt.Printf("%-4s %-25s %-25s %-10s %-10s\n", "DIR", "HOST", "LOCALSOCKET", "UPTIME(s)", "LASTDLV(s)") for _, c := range snap.TCPConns { - fmt.Printf("%-4s %-25s %-10s %-10s\n", - c.Dir, c.Host, ptrStr(c.UptimeSeconds), ptrStr(c.LastDlvSeconds)) + fmt.Printf("%-4s %-25s %-25s %-10s %-10s\n", + c.Dir, c.Host, c.LocalSocket, ptrStr(c.UptimeSeconds), ptrStr(c.LastDlvSeconds)) } - fmt.Println("\nsockets (host:port -> lastrcv/lastsnd s):") + fmt.Println("\nsockets by peer (host:port -> lastrcv/lastsnd s):") for host, s := range snap.Sockets { fmt.Printf(" %-25s lastrcv=%.1fs lastsnd=%.1fs\n", host, float64(s.LastRcvMs)/1000, float64(s.LastSndMs)/1000) } + + fmt.Println("\nsockets by local (host:port -> lastrcv/lastsnd s):") + for host, s := range snap.SocketsByLocal { + fmt.Printf(" %-25s lastrcv=%.1fs lastsnd=%.1fs\n", host, float64(s.LastRcvMs)/1000, float64(s.LastSndMs)/1000) + } } func ptrStr(p *int) string { From 2c79b1a0120db357b4f4d1a20c32522b6194a049 Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Wed, 15 Jul 2026 13:47:29 -0400 Subject: [PATCH 3/9] Add skupper debug sweep to detect and force-close idle TCP adaptor connections --- internal/cmd/skupper/common/flags.go | 7 + internal/cmd/skupper/debug/debug.go | 37 +++ .../cmd/skupper/debug/kube/conn_sweeper.go | 179 +++++++++++ .../cmd/skupper/debug/nonkube/conn_sweeper.go | 128 ++++++++ internal/cmd/skupper/debug/sweeper/README.md | 284 ------------------ .../cmd/skupper/debug/sweeper/criteria.go | 56 ++++ internal/cmd/skupper/debug/sweeper/gather.go | 29 +- .../skupper/debug/sweeper/gatherdump/main.go | 46 --- .../cmd/skupper/debug/sweeper/inetdiag.go | 86 ++++++ internal/cmd/skupper/debug/sweeper/kill.go | 35 +++ internal/cmd/skupper/debug/sweeper/sweeper.go | 95 ++++++ 11 files changed, 641 insertions(+), 341 deletions(-) create mode 100644 internal/cmd/skupper/debug/kube/conn_sweeper.go create mode 100644 internal/cmd/skupper/debug/nonkube/conn_sweeper.go delete mode 100644 internal/cmd/skupper/debug/sweeper/README.md create mode 100644 internal/cmd/skupper/debug/sweeper/criteria.go delete mode 100644 internal/cmd/skupper/debug/sweeper/gatherdump/main.go create mode 100644 internal/cmd/skupper/debug/sweeper/inetdiag.go create mode 100644 internal/cmd/skupper/debug/sweeper/kill.go create mode 100644 internal/cmd/skupper/debug/sweeper/sweeper.go diff --git a/internal/cmd/skupper/common/flags.go b/internal/cmd/skupper/common/flags.go index fa9487fae..6a1edc0fc 100644 --- a/internal/cmd/skupper/common/flags.go +++ b/internal/cmd/skupper/common/flags.go @@ -252,6 +252,13 @@ type CommandVersionFlags struct { type CommandDebugFlags struct { } +type CommandConnSweeperFlags struct { + URL string + IdleThreshold int + DryRun bool + Skmanage string +} + type CommandSystemUninstallFlags struct { Force bool } diff --git a/internal/cmd/skupper/debug/debug.go b/internal/cmd/skupper/debug/debug.go index 146192041..971d0b2fc 100644 --- a/internal/cmd/skupper/debug/debug.go +++ b/internal/cmd/skupper/debug/debug.go @@ -4,6 +4,7 @@ import ( "github.com/skupperproject/skupper/internal/cmd/skupper/common" "github.com/skupperproject/skupper/internal/cmd/skupper/debug/kube" "github.com/skupperproject/skupper/internal/cmd/skupper/debug/nonkube" + "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" "github.com/skupperproject/skupper/internal/config" "github.com/spf13/cobra" @@ -18,6 +19,42 @@ func NewCmdDebug() *cobra.Command { } platform := common.Platform(config.GetPlatform()) cmd.AddCommand(CmdDebugDumpFactory(platform)) + cmd.AddCommand(CmdDebugSweepFactory(platform)) + + return cmd +} + +func CmdDebugSweepFactory(configuredPlatform common.Platform) *cobra.Command { + kubeCommand := kube.NewCmdConnSweeper() + nonKubeCommand := nonkube.NewCmdConnSweeper() + + cmdDesc := common.SkupperCmdDescription{ + Use: "sweep", + Short: "Detect and kill idle TCP adaptor connections", + Long: `Queries the router management API for TCP adaptor connections, identifies +connections that have been idle beyond the threshold, and force-closes them +via adminStatus=deleted.`, + Example: "skupper debug sweep --url amqp://127.0.0.1:5672 --idle-threshold 14400", + } + + cmd := common.ConfigureCobraCommand(configuredPlatform, cmdDesc, kubeCommand, nonKubeCommand) + cmd.Hidden = true + + cmdFlags := common.CommandConnSweeperFlags{ + URL: sweeper.DefaultURL, + IdleThreshold: sweeper.DefaultIdleThreshold, + Skmanage: sweeper.DefaultSkmanage, + } + + cmd.Flags().StringVar(&cmdFlags.URL, "url", sweeper.DefaultURL, "Router management URL") + cmd.Flags().IntVar(&cmdFlags.IdleThreshold, "idle-threshold", sweeper.DefaultIdleThreshold, "Seconds with no data received before a connection is flagged as orphaned") + cmd.Flags().BoolVar(&cmdFlags.DryRun, "dry-run", false, "List idle connections without killing them") + cmd.Flags().StringVar(&cmdFlags.Skmanage, "skmanage", sweeper.DefaultSkmanage, "Path to the skmanage binary") + + kubeCommand.CobraCmd = cmd + kubeCommand.Flags = &cmdFlags + nonKubeCommand.CobraCmd = cmd + nonKubeCommand.Flags = &cmdFlags return cmd } diff --git a/internal/cmd/skupper/debug/kube/conn_sweeper.go b/internal/cmd/skupper/debug/kube/conn_sweeper.go new file mode 100644 index 000000000..97594b6f8 --- /dev/null +++ b/internal/cmd/skupper/debug/kube/conn_sweeper.go @@ -0,0 +1,179 @@ +package kube + +import ( + "context" + "fmt" + "time" + + "github.com/skupperproject/skupper/internal/cmd/skupper/common" + "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" + "github.com/skupperproject/skupper/internal/kube/client" + "github.com/spf13/cobra" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/kubernetes/scheme" + restclient "k8s.io/client-go/rest" + "k8s.io/client-go/tools/clientcmd" +) + +const ( + routerPodSelector = "app.kubernetes.io/name=skupper-router" + routerContainer = "router" + podExecTimeout = 30 * time.Second +) + +// CmdConnSweeper is the kubernetes entry point for `skupper debug sweep`. It +// finds every ready router pod and runs the sweeper against each one, with +// all commands exec'd inside the router container so they see the pod's own +// network namespace (no port-forward needed). +type CmdConnSweeper struct { + CobraCmd *cobra.Command + Flags *common.CommandConnSweeperFlags + KubeClient kubernetes.Interface + Rest *restclient.Config + Namespace string +} + +func NewCmdConnSweeper() *CmdConnSweeper { + return &CmdConnSweeper{} +} + +func (cmd *CmdConnSweeper) NewClient(cobraCommand *cobra.Command, args []string) { + cli, err := client.NewClient(cobraCommand.Flag("namespace").Value.String(), cobraCommand.Flag("context").Value.String(), cobraCommand.Flag("kubeconfig").Value.String()) + if err != nil { + return + } + cmd.KubeClient = cli.GetKubeClient() + cmd.Namespace = cli.Namespace + + loadingRules := clientcmd.NewDefaultClientConfigLoadingRules() + if kubeConfigPath := cobraCommand.Flag("kubeconfig").Value.String(); kubeConfigPath != "" { + loadingRules = &clientcmd.ClientConfigLoadingRules{ExplicitPath: kubeConfigPath} + } + kubeconfig := clientcmd.NewNonInteractiveDeferredLoadingClientConfig( + loadingRules, + &clientcmd.ConfigOverrides{CurrentContext: cobraCommand.Flag("context").Value.String()}, + ) + restconfig, err := kubeconfig.ClientConfig() + if err != nil { + return + } + restconfig.APIPath = "/api" + restconfig.GroupVersion = &corev1.SchemeGroupVersion + restconfig.NegotiatedSerializer = scheme.Codecs.WithoutConversion() + cmd.Rest = restconfig +} + +func (cmd *CmdConnSweeper) ValidateInput(args []string) error { + if cmd.Flags.IdleThreshold <= 0 { + return fmt.Errorf("--idle-threshold must be a positive number of seconds") + } + if cmd.KubeClient == nil || cmd.Rest == nil { + return fmt.Errorf("could not initialize kubernetes client") + } + return nil +} + +func (cmd *CmdConnSweeper) InputToOptions() {} + +func (cmd *CmdConnSweeper) Run() error { + podNames, err := cmd.findRouterPods() + if err != nil { + return err + } + + // Each ready replica has its own connections, so sweep every pod. A + // failure on one pod (e.g. exec cut off mid-sweep) must not leave the + // remaining replicas unswept. + var total sweeper.Result + var failedPods []string + for _, podName := range podNames { + fmt.Printf("=== router pod %s (namespace %s) ===\n", podName, cmd.Namespace) + res, err := sweeper.Run(sweeper.Config{ + URL: cmd.Flags.URL, + Skmanage: cmd.Flags.Skmanage, + IdleThresholdSecs: cmd.Flags.IdleThreshold, + DryRun: cmd.Flags.DryRun, + Exec: cmd.podExecer(podName), + }) + if err != nil { + fmt.Printf("sweep of pod %s failed: %v\n", podName, err) + failedPods = append(failedPods, podName) + continue + } + total.Total += res.Total + total.Killed += res.Killed + total.Skipped += res.Skipped + total.Failed += res.Failed + } + + if len(podNames) > 1 { + fmt.Printf("=== all pods: total:%d killed:%d skipped:%d failed:%d ===\n", + total.Total, total.Killed, total.Skipped, total.Failed) + } + if len(failedPods) > 0 { + return fmt.Errorf("sweep failed on %d of %d router pod(s): %v", len(failedPods), len(podNames), failedPods) + } + return nil +} + +func (cmd *CmdConnSweeper) WaitUntil() error { return nil } + +// findRouterPods returns every router pod whose router container is ready, +// with HA there are multiple replicas and each holds its own connections. +func (cmd *CmdConnSweeper) findRouterPods() ([]string, error) { + pods, err := cmd.KubeClient.CoreV1().Pods(cmd.Namespace).List(context.TODO(), metav1.ListOptions{LabelSelector: routerPodSelector}) + if err != nil { + return nil, fmt.Errorf("could not list router pods: %w", err) + } + var ready []string + for _, pod := range pods.Items { + // Phase stays "Running" even while a container crash-loops, so + // require the router container itself to be ready. + for _, cs := range pod.Status.ContainerStatuses { + if cs.Name == routerContainer && cs.Ready { + ready = append(ready, pod.Name) + } + } + } + if len(ready) == 0 { + return nil, fmt.Errorf("no ready skupper-router pod found in namespace %q", cmd.Namespace) + } + return ready, nil +} + +// podExecer returns a sweeper.Execer that runs argv inside the router +// container, where skmanage, python3 and the router's own network namespace +// are all available. +func (cmd *CmdConnSweeper) podExecer(podName string) sweeper.Execer { + return func(argv []string) ([]byte, error) { + // ExecCommandInContainer uses the deprecated context-less Stream, so + // it can't be cancelled directly; run it in a goroutine and give up + // on the sweep's behalf after podExecTimeout. The buffered channel + // lets the goroutine finish its send if Stream ever returns. + type execResult struct { + out []byte + err error + } + done := make(chan execResult, 1) + go func() { + out, err := client.ExecCommandInContainer(argv, podName, routerContainer, cmd.Namespace, cmd.KubeClient, cmd.Rest) + if err != nil { + done <- execResult{nil, err} + return + } + done <- execResult{out.Bytes(), nil} + }() + + select { + case r := <-done: + if r.err != nil { + return nil, fmt.Errorf("exec %q in pod %s failed: %w", argv[0], podName, r.err) + } + return r.out, nil + case <-time.After(podExecTimeout): + return nil, fmt.Errorf("exec %q in pod %s timed out after %s", argv[0], podName, podExecTimeout) + } + } +} diff --git a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go new file mode 100644 index 000000000..f817dd406 --- /dev/null +++ b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go @@ -0,0 +1,128 @@ +package nonkube + +import ( + "context" + "fmt" + "os/exec" + "time" + + "github.com/skupperproject/skupper/internal/cmd/skupper/common" + "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" + "github.com/skupperproject/skupper/internal/nonkube/client/runtime" + nonkubecommon "github.com/skupperproject/skupper/internal/nonkube/common" + "github.com/spf13/cobra" +) + +// Cert paths inside the router container (see compat.SiteStateRenderer, which +// mounts the runtime certs at /etc/skupper-router/runtime/certs). +const containerCertsPath = "/etc/skupper-router/runtime/certs/skupper-local-client" + +type CmdConnSweeper struct { + CobraCmd *cobra.Command + Flags *common.CommandConnSweeperFlags + namespace string + platform string + url string + sslArgs []string + exec sweeper.Execer +} + +func NewCmdConnSweeper() *CmdConnSweeper { + return &CmdConnSweeper{} +} + +func (cmd *CmdConnSweeper) NewClient(cobraCommand *cobra.Command, args []string) { + cmd.namespace = cobraCommand.Flag(common.FlagNameNamespace).Value.String() +} + +func (cmd *CmdConnSweeper) ValidateInput(args []string) error { + if cmd.Flags.IdleThreshold <= 0 { + return fmt.Errorf("--idle-threshold must be a positive number of seconds") + } + return nil +} + +func (cmd *CmdConnSweeper) InputToOptions() { + if cmd.namespace == "" { + cmd.namespace = "default" + } + cmd.url = cmd.Flags.URL + + // Detect the platform from the namespace's site config. + platformLoader := &nonkubecommon.NamespacePlatformLoader{} + platform, err := platformLoader.Load(cmd.namespace) + if err != nil { + return + } + cmd.platform = platform + + switch platform { + + case "podman", "docker": + // exec inside it so skmanage and the socket + // query see the router's network namespace. + containerName := cmd.namespace + "-skupper-router" + cmd.exec = containerExecer(platform, containerName) + if !cmd.urlOverridden() { + // The site's local management listener is amqps with client-cert + // auth; certs are mounted inside the container. + if url, err := runtime.GetLocalRouterAddress(cmd.namespace); err == nil { + cmd.url = url + cmd.sslArgs = sslArgs(containerCertsPath+"/tls.crt", containerCertsPath+"/tls.key", containerCertsPath+"/ca.crt") + } + } + default: // linux + if !cmd.urlOverridden() { + if url, err := runtime.GetLocalRouterAddress(cmd.namespace); err == nil { + certs := runtime.GetRuntimeTlsCert(cmd.namespace, "skupper-local-client") + cmd.url = url + cmd.sslArgs = sslArgs(certs.CertPath, certs.KeyPath, certs.CaPath) + } + } + } +} + +func (cmd *CmdConnSweeper) Run() error { + if cmd.exec != nil { + fmt.Printf("running against %s container %s-skupper-router\n", cmd.platform, cmd.namespace) + } + _, err := sweeper.Run(sweeper.Config{ + URL: cmd.url, + Skmanage: cmd.Flags.Skmanage, + IdleThresholdSecs: cmd.Flags.IdleThreshold, + DryRun: cmd.Flags.DryRun, + Exec: cmd.exec, + SkmanageExtraArgs: cmd.sslArgs, + }) + return err +} + +func (cmd *CmdConnSweeper) WaitUntil() error { return nil } + +// urlOverridden reports whether the user set --url. +func (cmd *CmdConnSweeper) urlOverridden() bool { + return cmd.CobraCmd != nil && cmd.CobraCmd.Flags().Changed("url") +} + +func sslArgs(cert, key, ca string) []string { + return []string{"--ssl-certificate", cert, "--ssl-key", key, "--ssl-trustfile", ca} +} + +// containerExecer returns a sweeper.Execer that runs argv in the router +// container via ` exec` so that skmanage and python3 are +// available. +func containerExecer(engine, containerName string) sweeper.Execer { + return func(argv []string) ([]byte, error) { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + full := append([]string{"exec", containerName}, argv...) + out, err := exec.CommandContext(ctx, engine, full...).Output() + if err != nil { + if ee, ok := err.(*exec.ExitError); ok && len(ee.Stderr) > 0 { + return nil, fmt.Errorf("%s exec %q failed: %w (stderr: %s)", engine, argv[0], err, ee.Stderr) + } + return nil, fmt.Errorf("%s exec %q failed: %w", engine, argv[0], err) + } + return out, nil + } +} diff --git a/internal/cmd/skupper/debug/sweeper/README.md b/internal/cmd/skupper/debug/sweeper/README.md deleted file mode 100644 index 3a7e82194..000000000 --- a/internal/cmd/skupper/debug/sweeper/README.md +++ /dev/null @@ -1,284 +0,0 @@ -# Connection sweeper — data gathering demo - -## Setup - -Needs: a built `skupper-router`, plus the `a.conf`, `b.conf`, and `random_traffic_sim.py` files below. - -Create/modify these three files in order to run the tests. - ---- - -`a.conf` — in `/build/` - -```conf -router { - mode: interior - id: A -} - -listener { - port: amqp - authenticatePeer: no - saslMechanisms: ANONYMOUS -} - -listener { - port: 10000 - authenticatePeer: no - saslMechanisms: ANONYMOUS - role: inter-router -} - -tcpListener { - port: 9090 - address: tcp-echo - name: tcp-echo-listener -} -``` - ---- - -`b.conf` — in `/build/` - -```conf -router { - mode: interior - id: B -} -listener { - host: 127.0.0.1 - port: 5673 - role: normal -} -connector { - host: 127.0.0.1 - port: 10000 - role: inter-router -} -tcpConnector { - address: tcp-echo - host: 127.0.0.1 - port: 9091 - name: tcp-echo-connector -} -``` - ---- - -`random_traffic_sim.py` — in `/scripts/` - -```python -#!/usr/bin/env python3 -""" -Opens random TCP connections to a router's TCP adaptor, each sending a small payload at a random interval. -""" - -import argparse -import random -import signal -import socket -import threading -import time -from datetime import datetime - - -def log(msg: str) -> None: - ts = datetime.now().strftime("%H:%M:%S") - print(f"[{ts}] {msg}", flush=True) - - -def _run_backend(port: int) -> None: - srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - srv.bind(("0.0.0.0", port)) - srv.listen(256) - while True: - try: - conn, _ = srv.accept() - threading.Thread(target=_backend_echo, args=(conn,), daemon=True).start() - except OSError: - break - - -def _backend_echo(conn: socket.socket) -> None: - try: - while True: - data = conn.recv(4096) - if not data: - return - conn.sendall(data) - except OSError: - pass - finally: - conn.close() - - -def _connection_worker(idx: int, host: str, port: int, min_interval: int, max_interval: int, - stop: threading.Event) -> None: - """Holds one TCP connection open for the whole run, sending a payload and - reading the echo back at an independently randomized interval.""" - sock = None - try: - sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - sock.connect((host, port)) - sock.sendall(f"conn-{idx} open\n".encode()) - sock.settimeout(5.0) - try: - sock.recv(4096) - except socket.timeout: - pass - - while not stop.is_set(): - interval = random.randint(min_interval, max_interval) - if stop.wait(interval): - break - try: - payload = f"conn-{idx} ping {int(time.time())}\n".encode() - sock.sendall(payload) - sock.recv(4096) - except OSError as e: - log(f"conn-{idx} error, exiting: {e}") - return - except OSError as e: - log(f"conn-{idx} failed to connect: {e}") - finally: - if sock: - sock.close() - - -def main() -> None: - parser = argparse.ArgumentParser( - prog="random_traffic_sim.py", - description="Steady TCP connections sending data at random intervals, for sweeper testing.", - formatter_class=argparse.RawDescriptionHelpFormatter, - epilog=__doc__, - ) - parser.add_argument("--host", default="127.0.0.1", help="Router host (default: 127.0.0.1)") - parser.add_argument("--port", type=int, default=9090, help="Router TCP adaptor port (default: 9090)") - parser.add_argument("--connections", type=int, default=20, help="Number of connections to open (default: 20)") - parser.add_argument("--min-interval", type=int, default=30, - help="Minimum seconds between sends on a connection (default: 30)") - parser.add_argument("--max-interval", type=int, default=500, - help="Maximum seconds between sends on a connection (default: 500)") - parser.add_argument("--backend-port", type=int, default=None, - help="Port for this script's own echo backend (default: --port + 1)") - args = parser.parse_args() - - backend_port = args.backend_port or (args.port + 1) - stop = threading.Event() - - def _on_sigint(_s, _f): - log("Interrupted — shutting down ...") - stop.set() - - signal.signal(signal.SIGINT, _on_sigint) - - threading.Thread(target=_run_backend, args=(backend_port,), daemon=True).start() - time.sleep(0.3) - log(f"Backend echoing on :{backend_port}") - - log(f"Opening {args.connections} connection(s) to {args.host}:{args.port}, " - f"each sending every {args.min_interval}-{args.max_interval}s ...") - threads = [] - for i in range(args.connections): - t = threading.Thread( - target=_connection_worker, - args=(i, args.host, args.port, args.min_interval, args.max_interval, stop), - daemon=True, - ) - t.start() - threads.append(t) - time.sleep(0.05) - - log("All connections open. Running until Ctrl+C.") - stop.wait() - - for t in threads: - t.join(timeout=2) - log("Exited.") - - -if __name__ == "__main__": - main() -``` - ---- - -Copy-paste, once per terminal (need 4 terminals): - -```bash -export SKUPPER_ROUTER=~/skupper-router # adjust if your checkout lives elsewhere -export SKUPPER_REPO=~/skupper -[ -f "$SKUPPER_ROUTER/build/config.sh" ] || echo "!! not found — fix SKUPPER_ROUTER above" -source "$SKUPPER_ROUTER/build/config.sh" -cd "$SKUPPER_REPO" -``` - -Structured as client -> Router A -> Router B -> echo backend - -## Run it - -**Terminal 1 — Router A:** -```bash -"$SKUPPER_ROUTER/build/router/skrouterd" -c "$SKUPPER_ROUTER/build/a.conf" -``` - -**Terminal 2 — Router B:** -```bash -"$SKUPPER_ROUTER/build/router/skrouterd" -c "$SKUPPER_ROUTER/build/b.conf" -``` - -**Terminal 3 — traffic** (starts its own echo backend on `:9091`): -```bash -python3 "$SKUPPER_ROUTER/scripts/random_traffic_sim.py" --connections 6 --min-interval 30 --max-interval 300 -``` - -Wait for `All connections open.`, then run the tests below in a fourth terminal. - ---- - -## Test 1 — the `in` side - -```bash -go run ./internal/cmd/skupper/debug/sweeper/gatherdump --url amqp://127.0.0.1:5672 -``` - -``` -DIR HOST LOCALSOCKET UPTIME(s) LASTDLV(s) -in 127.0.0.1:53484 127.0.0.1:9090 123 123 -in 127.0.0.1:53498 127.0.0.1:9090 123 123 -in 127.0.0.1:53506 127.0.0.1:9090 123 123 -in 127.0.0.1:53514 127.0.0.1:9090 123 123 -in 127.0.0.1:53528 127.0.0.1:9090 123 123 -in 127.0.0.1:53538 127.0.0.1:9090 123 123 -``` - -`HOST` is unique per connection. `LOCALSOCKET` is the same on every row cause it's the listener. - -The two `sockets by ...` sections below the table are the same kernel sockets indexed by each end — `in` connections are matched by their peer address, `out` connections by their local address. - - -## Test 2 — the `out` side - -Same command, Router B: - -```bash -go run ./internal/cmd/skupper/debug/sweeper/gatherdump --url amqp://127.0.0.1:5673 -``` - -``` -DIR HOST LOCALSOCKET UPTIME(s) LASTDLV(s) -out 127.0.0.1:9091 127.0.0.1:54944 133 133 -out 127.0.0.1:9091 127.0.0.1:54948 133 133 -out 127.0.0.1:9091 127.0.0.1:54956 133 133 -out 127.0.0.1:9091 127.0.0.1:54960 133 133 -out 127.0.0.1:9091 127.0.0.1:54972 133 133 -out 127.0.0.1:9091 127.0.0.1:54976 133 133 -``` - -## Clean up -```bash -pkill -f 'skrouterd -c .*a\.conf' -pkill -f 'skrouterd -c .*b\.conf' -pkill -f random_traffic_sim.py -``` diff --git a/internal/cmd/skupper/debug/sweeper/criteria.go b/internal/cmd/skupper/debug/sweeper/criteria.go new file mode 100644 index 000000000..752a91075 --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/criteria.go @@ -0,0 +1,56 @@ +package sweeper + +import ( + "fmt" + "time" +) + +// Decision is a connection flagged for closure +type Decision struct { + Conn connInfo + Reason string +} + +// Evaluate applies all kill criteria (currently only based off of time) to all +// connections and determines which connections are to be killed by kill.go. +func Evaluate(snap Snapshot, idleThreshold time.Duration) []Decision { + var orphans []Decision + for _, c := range snap.TCPConns { + if reason, isOrphan := orphanReason(c, snap, idleThreshold); isOrphan { + orphans = append(orphans, Decision{Conn: c, Reason: reason}) + } + } + return orphans +} + +// orphanReason decides whether connection, c, should be closed and why. When no socket can +// be matched the connection is skipped. +func orphanReason(c connInfo, snap Snapshot, idleThreshold time.Duration) (string, bool) { + sock, ok := matchSocket(c, snap) + if !ok { + return "", false + } + idle := time.Duration(min(sock.LastRcvMs, sock.LastSndMs)) * time.Millisecond + if idle >= idleThreshold { + return fmt.Sprintf("idle for %s", idle.Round(time.Second)), true + } + return "", false +} + +// matchSocket finds c's kernel socket. An 'in' connection's Host is its +// unique peer address; an 'out' connection's Host is the shared backend +// address, so it is matched by its unique LocalSocket instead. +func matchSocket(c connInfo, snap Snapshot) (socketInfo, bool) { + switch c.Dir { + case "in": + sock, ok := snap.Sockets[c.Host] + return sock, ok + case "out": + if c.LocalSocket == "" { + return socketInfo{}, false + } + sock, ok := snap.SocketsByLocal[c.LocalSocket] + return sock, ok + } + return socketInfo{}, false +} diff --git a/internal/cmd/skupper/debug/sweeper/gather.go b/internal/cmd/skupper/debug/sweeper/gather.go index b88fcb085..190292397 100644 --- a/internal/cmd/skupper/debug/sweeper/gather.go +++ b/internal/cmd/skupper/debug/sweeper/gather.go @@ -51,9 +51,9 @@ type Snapshot struct { SocketsByLocal map[string]socketInfo } -// Execer runs a command (argv) and returns its stdout. Both skmanage and the -// socket query go through it, so they always observe the same host — and so -// the same network namespace, which is what makes their results joinable. +// Execer runs a command (argv) and returns its stdout. LocalExec runs on +// this host; the kube variant execs inside the router pod instead, so both +// skmanage and the socket query see the router's network namespace. type Execer func(argv []string) ([]byte, error) func LocalExec(argv []string) ([]byte, error) { @@ -61,9 +61,11 @@ func LocalExec(argv []string) ([]byte, error) { } // Gather queries the router for its TCP adaptor connections and cross -// references them with kernel socket state. Discards non-TCP-adaptor connections -func Gather(execFn Execer, skmanageBin, url string) (Snapshot, error) { - raw, err := runSkmanage(execFn, skmanageBin, url, "QUERY", "--type="+ConnType) +// references them with kernel socket state. Discards non-TCP-adaptor connections. +// extraArgs are appended to the skmanage invocation (e.g. --ssl-certificate +// options when the management endpoint is amqps). +func Gather(execFn Execer, skmanageBin, url string, extraArgs ...string) (Snapshot, error) { + raw, err := runSkmanage(execFn, skmanageBin, url, extraArgs, "QUERY", "--type="+ConnType) if err != nil { return Snapshot{}, fmt.Errorf("could not query router at %s: %w", url, err) } @@ -93,14 +95,18 @@ func isTCPAdaptorConn(c connInfo) bool { return c.Container == tcpContainer && c.Host != egressDispatch } -// gatherSockets reads kernel socket state via `ss -tin`. A host without ss (to be patched later) -// yields no sockets, which leaves every connection unmatched and untouched. +// gatherSockets reads kernel socket state, preferring `ss -tin` and falling +// back to the python netlink script when ss isn't available (e.g. inside the +// router container, which ships python3 but not iproute). func gatherSockets(execFn Execer) (byPeer, byLocal map[string]socketInfo) { - out, err := execFn([]string{"ss", "-tin"}) + if out, err := execFn([]string{"ss", "-tin"}); err == nil { + return socketsFromSS(out) + } + out, err := execFn([]string{"python3", "-c", inetDiagScript}) if err != nil { return map[string]socketInfo{}, map[string]socketInfo{} } - return socketsFromSS(out) + return socketsFromDiagOutput(out) } // socketsFromSS builds two {lastrcv, lastsnd} maps — one keyed by peer @@ -158,8 +164,9 @@ func extractMsField(line, key string) int { return val } -func runSkmanage(execFn Execer, bin, url string, args ...string) ([]byte, error) { +func runSkmanage(execFn Execer, bin, url string, extraArgs []string, args ...string) ([]byte, error) { argv := append([]string{bin, "--bus", url}, args...) + argv = append(argv, extraArgs...) out, err := execFn(argv) if err != nil { return nil, fmt.Errorf("skmanage failed: %w", err) diff --git a/internal/cmd/skupper/debug/sweeper/gatherdump/main.go b/internal/cmd/skupper/debug/sweeper/gatherdump/main.go deleted file mode 100644 index a3115c554..000000000 --- a/internal/cmd/skupper/debug/sweeper/gatherdump/main.go +++ /dev/null @@ -1,46 +0,0 @@ -// Throwaway debug tool: calls sweeper.Gather directly and dumps the results for debugging -package main - -import ( - "flag" - "fmt" - "os" - - "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" -) - -func main() { - url := flag.String("url", "amqp://127.0.0.1:5672", "Router management URL") - skmanageBin := flag.String("skmanage", "skmanage", "Path to skmanage binary") - flag.Parse() - - snap, err := sweeper.Gather(sweeper.LocalExec, *skmanageBin, *url) - if err != nil { - fmt.Fprintln(os.Stderr, "gather failed:", err) - os.Exit(1) - } - - fmt.Printf("gathered at %s from %s\n\n", snap.Now.Format("15:04:05"), *url) - fmt.Printf("%-4s %-25s %-25s %-10s %-10s\n", "DIR", "HOST", "LOCALSOCKET", "UPTIME(s)", "LASTDLV(s)") - for _, c := range snap.TCPConns { - fmt.Printf("%-4s %-25s %-25s %-10s %-10s\n", - c.Dir, c.Host, c.LocalSocket, ptrStr(c.UptimeSeconds), ptrStr(c.LastDlvSeconds)) - } - - fmt.Println("\nsockets by peer (host:port -> lastrcv/lastsnd s):") - for host, s := range snap.Sockets { - fmt.Printf(" %-25s lastrcv=%.1fs lastsnd=%.1fs\n", host, float64(s.LastRcvMs)/1000, float64(s.LastSndMs)/1000) - } - - fmt.Println("\nsockets by local (host:port -> lastrcv/lastsnd s):") - for host, s := range snap.SocketsByLocal { - fmt.Printf(" %-25s lastrcv=%.1fs lastsnd=%.1fs\n", host, float64(s.LastRcvMs)/1000, float64(s.LastSndMs)/1000) - } -} - -func ptrStr(p *int) string { - if p == nil { - return "-" - } - return fmt.Sprintf("%d", *p) -} diff --git a/internal/cmd/skupper/debug/sweeper/inetdiag.go b/internal/cmd/skupper/debug/sweeper/inetdiag.go new file mode 100644 index 000000000..ac3fd298b --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/inetdiag.go @@ -0,0 +1,86 @@ +package sweeper + +import ( + "strconv" + "strings" +) + +// inetDiagScript is a python stand-in for `ss -tin`, used when ss is +// not available like when the router image ships python3 but no iproute2. It has to +// run where the router runs (sockets are per network namespace), so it can +// only use what the router image provides. +// +// Like ss, it queries the kernel over NETLINK_SOCK_DIAG, but reports only +// what the sweeper needs: established IPv4 TCP sockets (Script might need to be updated to cover IPv6), one per line, with +// the two tcp_info idle timers: +// +// + +const inetDiagScript = ` +import socket, struct + +NETLINK_SOCK_DIAG = 4 +SOCK_DIAG_BY_FAMILY = 20 +NLM_F_REQUEST_DUMP = 0x301 +NLMSG_DONE = 3 +NLMSG_ERROR = 2 +INET_DIAG_INFO = 2 + +s = socket.socket(socket.AF_NETLINK, socket.SOCK_RAW, NETLINK_SOCK_DIAG) +# inet_diag_req_v2: family, protocol, ext (request tcp_info), pad, +# states bitmask (established only), zeroed socket id +req = struct.pack("=BBBBI48s", socket.AF_INET, socket.IPPROTO_TCP, + 1 << (INET_DIAG_INFO - 1), 0, 1 << 1, b"") +s.send(struct.pack("=IHHII", 16 + len(req), SOCK_DIAG_BY_FAMILY, + NLM_F_REQUEST_DUMP, 1, 0) + req) + +done = False +while not done: + data = s.recv(1 << 20) + off = 0 + while off + 16 <= len(data): + ln, typ = struct.unpack_from("=IH", data, off) + if ln < 16 or typ in (NLMSG_DONE, NLMSG_ERROR): + done = True + break + # inet_diag_msg: sport/dport are big-endian at +4/+6, src at +8, + # dst at +24 (IPv4 uses the first 4 of 16 address bytes) + sport, dport = struct.unpack_from(">HH", data, off + 20) + src = socket.inet_ntoa(data[off + 24:off + 28]) + dst = socket.inet_ntoa(data[off + 40:off + 44]) + lastrcv = lastsnd = 0 + aoff = off + 16 + 72 # rtattrs follow the 72-byte inet_diag_msg + while aoff + 4 <= off + ln: + alen, atype = struct.unpack_from("=HH", data, aoff) + if alen < 4: + break + if atype == INET_DIAG_INFO and alen >= 4 + 56: + # tcp_info: tcpi_last_data_sent at +44, tcpi_last_data_recv at +52 + lastsnd = struct.unpack_from("=I", data, aoff + 4 + 44)[0] + lastrcv = struct.unpack_from("=I", data, aoff + 4 + 52)[0] + aoff += (alen + 3) & ~3 + print("%s:%d %s:%d %d %d" % (src, sport, dst, dport, lastrcv, lastsnd)) + off += (ln + 3) & ~3 +` + +// socketsFromDiagOutput parses inetDiagScript's output into the same two +// maps socketsFromSS produces. +func socketsFromDiagOutput(out []byte) (byPeer, byLocal map[string]socketInfo) { + byPeer = map[string]socketInfo{} + byLocal = map[string]socketInfo{} + for _, line := range strings.Split(string(out), "\n") { + f := strings.Fields(line) + if len(f) != 4 { + continue + } + rcv, err1 := strconv.Atoi(f[2]) + snd, err2 := strconv.Atoi(f[3]) + if err1 != nil || err2 != nil { + continue + } + sock := socketInfo{LastRcvMs: rcv, LastSndMs: snd} + byLocal[f[0]] = sock + byPeer[f[1]] = sock + } + return byPeer, byLocal +} diff --git a/internal/cmd/skupper/debug/sweeper/kill.go b/internal/cmd/skupper/debug/sweeper/kill.go new file mode 100644 index 000000000..f1dfbac94 --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/kill.go @@ -0,0 +1,35 @@ +package sweeper + +// killAll force-closes every TCP connection that criteria.go decided sequentially +func killAll(execFn Execer, skmanageBin, url string, extraArgs []string, decisions []Decision) (killed, failed int) { + for _, d := range decisions { + _, err := runSkmanage(execFn, skmanageBin, url, extraArgs, + "UPDATE", + "--type="+ConnType, + "--identity="+d.Conn.Identity, + "adminStatus=deleted", + ) + if err == nil { + logf(" id=%s host=%s dir=%s uptime=%s reason=%s → killed", + d.Conn.Identity, d.Conn.Host, d.Conn.Dir, + fmtSeconds(d.Conn.UptimeSeconds), d.Reason) + killed++ + continue + } + // Killing one half of a proxied pair cascade-closes the other half, + // so a later kill of that half fails with "not found". Confirm the + // connection is really gone before calling it a failure. + if _, readErr := runSkmanage(execFn, skmanageBin, url, extraArgs, + "READ", "--type="+ConnType, "--identity="+d.Conn.Identity); readErr != nil { + logf(" id=%s host=%s dir=%s reason=%s → already closed", + d.Conn.Identity, d.Conn.Host, d.Conn.Dir, d.Reason) + killed++ + } else { + logf(" id=%s host=%s dir=%s reason=%s → failed: %s", + d.Conn.Identity, d.Conn.Host, d.Conn.Dir, + d.Reason, err.Error()) + failed++ + } + } + return killed, failed +} diff --git a/internal/cmd/skupper/debug/sweeper/sweeper.go b/internal/cmd/skupper/debug/sweeper/sweeper.go new file mode 100644 index 000000000..759123c05 --- /dev/null +++ b/internal/cmd/skupper/debug/sweeper/sweeper.go @@ -0,0 +1,95 @@ +package sweeper + +import ( + "fmt" + "time" +) + +const ( + DefaultURL = "amqp://127.0.0.1:5672" + DefaultIdleThreshold = 4 * 3600 + DefaultSkmanage = "skmanage" +) + +type Config struct { + URL string + Skmanage string + IdleThresholdSecs int + DryRun bool + // Exec runs skmanage and the socket query. Nil means locally; the kube + // variant supplies a pod-exec so both run inside the router container. + Exec Execer + // SkmanageExtraArgs is appended to every skmanage invocation — e.g. + // --ssl-certificate/--ssl-key/--ssl-trustfile when the management + // endpoint is amqps (nonkube sites). + SkmanageExtraArgs []string +} + +type Result struct { + Total int + Killed int + Skipped int + Failed int +} + +// Run ties the two independent stages together: Gather (gather.go) collects +// raw router + kernel state, Evaluate (criteria.go) applies the idle-time +// criteria against that state, and killAll (kill.go) carries out whatever +// Evaluate decided. +func Run(cfg Config) (Result, error) { + if cfg.Exec == nil { + cfg.Exec = LocalExec + } + snap, err := Gather(cfg.Exec, cfg.Skmanage, cfg.URL, cfg.SkmanageExtraArgs...) + if err != nil { + return Result{}, err + } + + toKill := Evaluate(snap, time.Duration(cfg.IdleThresholdSecs)*time.Second) + logf("total:%d idle-orphan:%d", len(snap.TCPConns), len(toKill)) + + if len(toKill) == 0 { + logf("No idle/orphaned connections found.") + return Result{Total: len(snap.TCPConns)}, nil + } + + if cfg.DryRun { + logf("DRY RUN — would kill %d connection(s):", len(toKill)) + for _, d := range toKill { + fmt.Printf(" id=%-6s host=%-25s dir=%s uptime=%-10s reason=%s\n", + d.Conn.Identity, d.Conn.Host, d.Conn.Dir, fmtSeconds(d.Conn.UptimeSeconds), d.Reason) + } + return Result{Total: len(snap.TCPConns), Skipped: len(toKill)}, nil + } + + logf("--- KILLING %d connection(s) ---", len(toKill)) + killed, failed := killAll(cfg.Exec, cfg.Skmanage, cfg.URL, cfg.SkmanageExtraArgs, toKill) + + return Result{Total: len(snap.TCPConns), Killed: killed, Failed: failed}, nil +} + +func logf(format string, args ...any) { + ts := time.Now().Format("15:04:05") + fmt.Printf("["+ts+"] "+format+"\n", args...) +} + +func fmtSeconds(s *int) string { + if s == nil { + return "never" + } + return fmtDuration(time.Duration(*s) * time.Second) +} + +func fmtDuration(d time.Duration) string { + sec := int(d.Seconds()) + h := sec / 3600 + m := (sec % 3600) / 60 + s := sec % 60 + if h > 0 { + return fmt.Sprintf("%dh%02dm%02ds", h, m, s) + } + if m > 0 { + return fmt.Sprintf("%dm%02ds", m, s) + } + return fmt.Sprintf("%ds", s) +} From f529c33e3871ec1169131f3d506b1b5c185f9f07 Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Wed, 15 Jul 2026 14:52:35 -0400 Subject: [PATCH 4/9] addressed coderabbit concerns --- .../cmd/skupper/debug/kube/conn_sweeper.go | 29 ++++++++++++++----- .../cmd/skupper/debug/nonkube/conn_sweeper.go | 3 +- 2 files changed, 23 insertions(+), 9 deletions(-) diff --git a/internal/cmd/skupper/debug/kube/conn_sweeper.go b/internal/cmd/skupper/debug/kube/conn_sweeper.go index 97594b6f8..c788d5c1f 100644 --- a/internal/cmd/skupper/debug/kube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/kube/conn_sweeper.go @@ -33,6 +33,7 @@ type CmdConnSweeper struct { KubeClient kubernetes.Interface Rest *restclient.Config Namespace string + clientErr error } func NewCmdConnSweeper() *CmdConnSweeper { @@ -40,35 +41,47 @@ func NewCmdConnSweeper() *CmdConnSweeper { } func (cmd *CmdConnSweeper) NewClient(cobraCommand *cobra.Command, args []string) { - cli, err := client.NewClient(cobraCommand.Flag("namespace").Value.String(), cobraCommand.Flag("context").Value.String(), cobraCommand.Flag("kubeconfig").Value.String()) + cmd.CobraCmd = cobraCommand + namespaceFlag, _ := cobraCommand.Flags().GetString("namespace") + contextFlag, _ := cobraCommand.Flags().GetString("context") + kubeconfigFlag, _ := cobraCommand.Flags().GetString("kubeconfig") + + cli, err := client.NewClient(namespaceFlag, contextFlag, kubeconfigFlag) if err != nil { + cmd.clientErr = fmt.Errorf("failed to initialize kubernetes client: %w", err) return } cmd.KubeClient = cli.GetKubeClient() cmd.Namespace = cli.Namespace loadingRules := clientcmd.NewDefaultClientConfigLoadingRules() - if kubeConfigPath := cobraCommand.Flag("kubeconfig").Value.String(); kubeConfigPath != "" { - loadingRules = &clientcmd.ClientConfigLoadingRules{ExplicitPath: kubeConfigPath} + if kubeconfigFlag != "" { + loadingRules = &clientcmd.ClientConfigLoadingRules{ExplicitPath: kubeconfigFlag} } kubeconfig := clientcmd.NewNonInteractiveDeferredLoadingClientConfig( loadingRules, - &clientcmd.ConfigOverrides{CurrentContext: cobraCommand.Flag("context").Value.String()}, + &clientcmd.ConfigOverrides{CurrentContext: contextFlag}, ) restconfig, err := kubeconfig.ClientConfig() if err != nil { + cmd.clientErr = fmt.Errorf("failed to build kubernetes client config: %w", err) return } - restconfig.APIPath = "/api" - restconfig.GroupVersion = &corev1.SchemeGroupVersion - restconfig.NegotiatedSerializer = scheme.Codecs.WithoutConversion() - cmd.Rest = restconfig + + execConfig := *restconfig + execConfig.APIPath = "/api" + execConfig.GroupVersion = &corev1.SchemeGroupVersion + execConfig.NegotiatedSerializer = scheme.Codecs.WithoutConversion() + cmd.Rest = &execConfig } func (cmd *CmdConnSweeper) ValidateInput(args []string) error { if cmd.Flags.IdleThreshold <= 0 { return fmt.Errorf("--idle-threshold must be a positive number of seconds") } + if cmd.clientErr != nil { + return cmd.clientErr + } if cmd.KubeClient == nil || cmd.Rest == nil { return fmt.Errorf("could not initialize kubernetes client") } diff --git a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go index f817dd406..1faa46afd 100644 --- a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go @@ -32,7 +32,8 @@ func NewCmdConnSweeper() *CmdConnSweeper { } func (cmd *CmdConnSweeper) NewClient(cobraCommand *cobra.Command, args []string) { - cmd.namespace = cobraCommand.Flag(common.FlagNameNamespace).Value.String() + cmd.CobraCmd = cobraCommand + cmd.namespace, _ = cobraCommand.Flags().GetString(common.FlagNameNamespace) } func (cmd *CmdConnSweeper) ValidateInput(args []string) error { From 850c7387219758561599df266843253095440c95 Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Mon, 20 Jul 2026 11:22:13 -0400 Subject: [PATCH 5/9] Addressed Noe + Andy's Comments --- internal/cmd/skupper/common/flags.go | 2 - internal/cmd/skupper/debug/debug.go | 10 +-- .../cmd/skupper/debug/kube/conn_sweeper.go | 39 +++++------ .../cmd/skupper/debug/nonkube/conn_sweeper.go | 64 +++++++++---------- internal/cmd/skupper/debug/sweeper/gather.go | 2 +- .../cmd/skupper/debug/sweeper/inetdiag.go | 5 +- internal/cmd/skupper/debug/sweeper/kill.go | 3 +- internal/cmd/skupper/debug/sweeper/sweeper.go | 7 +- 8 files changed, 57 insertions(+), 75 deletions(-) diff --git a/internal/cmd/skupper/common/flags.go b/internal/cmd/skupper/common/flags.go index 6a1edc0fc..5681e56d4 100644 --- a/internal/cmd/skupper/common/flags.go +++ b/internal/cmd/skupper/common/flags.go @@ -253,10 +253,8 @@ type CommandDebugFlags struct { } type CommandConnSweeperFlags struct { - URL string IdleThreshold int DryRun bool - Skmanage string } type CommandSystemUninstallFlags struct { diff --git a/internal/cmd/skupper/debug/debug.go b/internal/cmd/skupper/debug/debug.go index 971d0b2fc..2dc90c8db 100644 --- a/internal/cmd/skupper/debug/debug.go +++ b/internal/cmd/skupper/debug/debug.go @@ -34,22 +34,16 @@ func CmdDebugSweepFactory(configuredPlatform common.Platform) *cobra.Command { Long: `Queries the router management API for TCP adaptor connections, identifies connections that have been idle beyond the threshold, and force-closes them via adminStatus=deleted.`, - Example: "skupper debug sweep --url amqp://127.0.0.1:5672 --idle-threshold 14400", + Example: "skupper debug sweep --idle-threshold 14400", } cmd := common.ConfigureCobraCommand(configuredPlatform, cmdDesc, kubeCommand, nonKubeCommand) cmd.Hidden = true - cmdFlags := common.CommandConnSweeperFlags{ - URL: sweeper.DefaultURL, - IdleThreshold: sweeper.DefaultIdleThreshold, - Skmanage: sweeper.DefaultSkmanage, - } + var cmdFlags common.CommandConnSweeperFlags - cmd.Flags().StringVar(&cmdFlags.URL, "url", sweeper.DefaultURL, "Router management URL") cmd.Flags().IntVar(&cmdFlags.IdleThreshold, "idle-threshold", sweeper.DefaultIdleThreshold, "Seconds with no data received before a connection is flagged as orphaned") cmd.Flags().BoolVar(&cmdFlags.DryRun, "dry-run", false, "List idle connections without killing them") - cmd.Flags().StringVar(&cmdFlags.Skmanage, "skmanage", sweeper.DefaultSkmanage, "Path to the skmanage binary") kubeCommand.CobraCmd = cmd kubeCommand.Flags = &cmdFlags diff --git a/internal/cmd/skupper/debug/kube/conn_sweeper.go b/internal/cmd/skupper/debug/kube/conn_sweeper.go index c788d5c1f..2561010d0 100644 --- a/internal/cmd/skupper/debug/kube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/kube/conn_sweeper.go @@ -2,12 +2,14 @@ package kube import ( "context" + "errors" "fmt" "time" "github.com/skupperproject/skupper/internal/cmd/skupper/common" "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" "github.com/skupperproject/skupper/internal/kube/client" + "github.com/skupperproject/skupper/internal/utils/validator" "github.com/spf13/cobra" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -76,36 +78,38 @@ func (cmd *CmdConnSweeper) NewClient(cobraCommand *cobra.Command, args []string) } func (cmd *CmdConnSweeper) ValidateInput(args []string) error { - if cmd.Flags.IdleThreshold <= 0 { - return fmt.Errorf("--idle-threshold must be a positive number of seconds") + var validationErrors []error + numberValidator := validator.NewNumberValidator() + numberValidator.IncludeZero = false + if ok, err := numberValidator.Evaluate(cmd.Flags.IdleThreshold); !ok { + validationErrors = append(validationErrors, fmt.Errorf("idle-threshold is not valid: %s", err)) } + return errors.Join(validationErrors...) +} + +func (cmd *CmdConnSweeper) InputToOptions() {} + +func (cmd *CmdConnSweeper) Run() error { if cmd.clientErr != nil { return cmd.clientErr } if cmd.KubeClient == nil || cmd.Rest == nil { return fmt.Errorf("could not initialize kubernetes client") } - return nil -} -func (cmd *CmdConnSweeper) InputToOptions() {} - -func (cmd *CmdConnSweeper) Run() error { podNames, err := cmd.findRouterPods() if err != nil { return err } - // Each ready replica has its own connections, so sweep every pod. A - // failure on one pod (e.g. exec cut off mid-sweep) must not leave the - // remaining replicas unswept. + // Each ready replica has its own connections, so sweep every pod. var total sweeper.Result var failedPods []string for _, podName := range podNames { fmt.Printf("=== router pod %s (namespace %s) ===\n", podName, cmd.Namespace) res, err := sweeper.Run(sweeper.Config{ - URL: cmd.Flags.URL, - Skmanage: cmd.Flags.Skmanage, + URL: sweeper.DefaultURL, + Skmanage: sweeper.DefaultSkmanage, IdleThresholdSecs: cmd.Flags.IdleThreshold, DryRun: cmd.Flags.DryRun, Exec: cmd.podExecer(podName), @@ -133,8 +137,6 @@ func (cmd *CmdConnSweeper) Run() error { func (cmd *CmdConnSweeper) WaitUntil() error { return nil } -// findRouterPods returns every router pod whose router container is ready, -// with HA there are multiple replicas and each holds its own connections. func (cmd *CmdConnSweeper) findRouterPods() ([]string, error) { pods, err := cmd.KubeClient.CoreV1().Pods(cmd.Namespace).List(context.TODO(), metav1.ListOptions{LabelSelector: routerPodSelector}) if err != nil { @@ -142,8 +144,6 @@ func (cmd *CmdConnSweeper) findRouterPods() ([]string, error) { } var ready []string for _, pod := range pods.Items { - // Phase stays "Running" even while a container crash-loops, so - // require the router container itself to be ready. for _, cs := range pod.Status.ContainerStatuses { if cs.Name == routerContainer && cs.Ready { ready = append(ready, pod.Name) @@ -157,14 +157,9 @@ func (cmd *CmdConnSweeper) findRouterPods() ([]string, error) { } // podExecer returns a sweeper.Execer that runs argv inside the router -// container, where skmanage, python3 and the router's own network namespace -// are all available. +// container. func (cmd *CmdConnSweeper) podExecer(podName string) sweeper.Execer { return func(argv []string) ([]byte, error) { - // ExecCommandInContainer uses the deprecated context-less Stream, so - // it can't be cancelled directly; run it in a goroutine and give up - // on the sweep's behalf after podExecTimeout. The buffered channel - // lets the goroutine finish its send if Stream ever returns. type execResult struct { out []byte err error diff --git a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go index 1faa46afd..2cd9a35d5 100644 --- a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go @@ -2,6 +2,7 @@ package nonkube import ( "context" + "errors" "fmt" "os/exec" "time" @@ -10,12 +11,15 @@ import ( "github.com/skupperproject/skupper/internal/cmd/skupper/debug/sweeper" "github.com/skupperproject/skupper/internal/nonkube/client/runtime" nonkubecommon "github.com/skupperproject/skupper/internal/nonkube/common" + "github.com/skupperproject/skupper/internal/utils/validator" "github.com/spf13/cobra" ) -// Cert paths inside the router container (see compat.SiteStateRenderer, which -// mounts the runtime certs at /etc/skupper-router/runtime/certs). const containerCertsPath = "/etc/skupper-router/runtime/certs/skupper-local-client" +const ( + containerSkmanage = "/bin/skmanage" + systemdSkmanage = "/usr/bin/skmanage" +) type CmdConnSweeper struct { CobraCmd *cobra.Command @@ -23,6 +27,7 @@ type CmdConnSweeper struct { namespace string platform string url string + skmanage string sslArgs []string exec sweeper.Execer } @@ -37,19 +42,20 @@ func (cmd *CmdConnSweeper) NewClient(cobraCommand *cobra.Command, args []string) } func (cmd *CmdConnSweeper) ValidateInput(args []string) error { - if cmd.Flags.IdleThreshold <= 0 { - return fmt.Errorf("--idle-threshold must be a positive number of seconds") + var validationErrors []error + numberValidator := validator.NewNumberValidator() + numberValidator.IncludeZero = false + if ok, err := numberValidator.Evaluate(cmd.Flags.IdleThreshold); !ok { + validationErrors = append(validationErrors, fmt.Errorf("idle-threshold is not valid: %s", err)) } - return nil + return errors.Join(validationErrors...) } func (cmd *CmdConnSweeper) InputToOptions() { if cmd.namespace == "" { cmd.namespace = "default" } - cmd.url = cmd.Flags.URL - // Detect the platform from the namespace's site config. platformLoader := &nonkubecommon.NamespacePlatformLoader{} platform, err := platformLoader.Load(cmd.namespace) if err != nil { @@ -57,39 +63,34 @@ func (cmd *CmdConnSweeper) InputToOptions() { } cmd.platform = platform - switch platform { + url, err := runtime.GetLocalRouterAddress(cmd.namespace) + if err != nil { + return + } + cmd.url = url + switch platform { case "podman", "docker": - // exec inside it so skmanage and the socket - // query see the router's network namespace. - containerName := cmd.namespace + "-skupper-router" - cmd.exec = containerExecer(platform, containerName) - if !cmd.urlOverridden() { - // The site's local management listener is amqps with client-cert - // auth; certs are mounted inside the container. - if url, err := runtime.GetLocalRouterAddress(cmd.namespace); err == nil { - cmd.url = url - cmd.sslArgs = sslArgs(containerCertsPath+"/tls.crt", containerCertsPath+"/tls.key", containerCertsPath+"/ca.crt") - } - } - default: // linux - if !cmd.urlOverridden() { - if url, err := runtime.GetLocalRouterAddress(cmd.namespace); err == nil { - certs := runtime.GetRuntimeTlsCert(cmd.namespace, "skupper-local-client") - cmd.url = url - cmd.sslArgs = sslArgs(certs.CertPath, certs.KeyPath, certs.CaPath) - } - } + cmd.skmanage = containerSkmanage + cmd.exec = containerExecer(platform, cmd.namespace+"-skupper-router") + cmd.sslArgs = sslArgs(containerCertsPath+"/tls.crt", containerCertsPath+"/tls.key", containerCertsPath+"/ca.crt") + default: + cmd.skmanage = systemdSkmanage + certs := runtime.GetRuntimeTlsCert(cmd.namespace, "skupper-local-client") + cmd.sslArgs = sslArgs(certs.CertPath, certs.KeyPath, certs.CaPath) } } func (cmd *CmdConnSweeper) Run() error { + if cmd.url == "" || cmd.skmanage == "" { + return fmt.Errorf("could not determine router management address for namespace %q", cmd.namespace) + } if cmd.exec != nil { fmt.Printf("running against %s container %s-skupper-router\n", cmd.platform, cmd.namespace) } _, err := sweeper.Run(sweeper.Config{ URL: cmd.url, - Skmanage: cmd.Flags.Skmanage, + Skmanage: cmd.skmanage, IdleThresholdSecs: cmd.Flags.IdleThreshold, DryRun: cmd.Flags.DryRun, Exec: cmd.exec, @@ -100,11 +101,6 @@ func (cmd *CmdConnSweeper) Run() error { func (cmd *CmdConnSweeper) WaitUntil() error { return nil } -// urlOverridden reports whether the user set --url. -func (cmd *CmdConnSweeper) urlOverridden() bool { - return cmd.CobraCmd != nil && cmd.CobraCmd.Flags().Changed("url") -} - func sslArgs(cert, key, ca string) []string { return []string{"--ssl-certificate", cert, "--ssl-key", key, "--ssl-trustfile", ca} } diff --git a/internal/cmd/skupper/debug/sweeper/gather.go b/internal/cmd/skupper/debug/sweeper/gather.go index 190292397..9fe9e77f6 100644 --- a/internal/cmd/skupper/debug/sweeper/gather.go +++ b/internal/cmd/skupper/debug/sweeper/gather.go @@ -31,7 +31,7 @@ type connInfo struct { } // socketInfo is the kernel's view of a TCP socket, as reported by `ss -tin`. -// LastRcvMs/LastSndMs come straight from TCP_INFO and should show actual activity vs lastDlvSeconds +// LastRcvMs/LastSndMs come from the kernel's TCP_INFO. type socketInfo struct { LastRcvMs int LastSndMs int diff --git a/internal/cmd/skupper/debug/sweeper/inetdiag.go b/internal/cmd/skupper/debug/sweeper/inetdiag.go index ac3fd298b..4498deaba 100644 --- a/internal/cmd/skupper/debug/sweeper/inetdiag.go +++ b/internal/cmd/skupper/debug/sweeper/inetdiag.go @@ -10,11 +10,10 @@ import ( // run where the router runs (sockets are per network namespace), so it can // only use what the router image provides. // -// Like ss, it queries the kernel over NETLINK_SOCK_DIAG, but reports only -// what the sweeper needs: established IPv4 TCP sockets (Script might need to be updated to cover IPv6), one per line, with -// the two tcp_info idle timers: // // +// +// TODO: IPv4 only; extend for IPv6. const inetDiagScript = ` import socket, struct diff --git a/internal/cmd/skupper/debug/sweeper/kill.go b/internal/cmd/skupper/debug/sweeper/kill.go index f1dfbac94..21d15be58 100644 --- a/internal/cmd/skupper/debug/sweeper/kill.go +++ b/internal/cmd/skupper/debug/sweeper/kill.go @@ -1,6 +1,7 @@ package sweeper -// killAll force-closes every TCP connection that criteria.go decided sequentially +// killAll force-closes each flagged connection by setting adminStatus=deleted +// via skmanage. func killAll(execFn Execer, skmanageBin, url string, extraArgs []string, decisions []Decision) (killed, failed int) { for _, d := range decisions { _, err := runSkmanage(execFn, skmanageBin, url, extraArgs, diff --git a/internal/cmd/skupper/debug/sweeper/sweeper.go b/internal/cmd/skupper/debug/sweeper/sweeper.go index 759123c05..824407dbc 100644 --- a/internal/cmd/skupper/debug/sweeper/sweeper.go +++ b/internal/cmd/skupper/debug/sweeper/sweeper.go @@ -7,8 +7,8 @@ import ( const ( DefaultURL = "amqp://127.0.0.1:5672" - DefaultIdleThreshold = 4 * 3600 DefaultSkmanage = "skmanage" + DefaultIdleThreshold = 4 * 3600 // 4 hours ) type Config struct { @@ -16,8 +16,7 @@ type Config struct { Skmanage string IdleThresholdSecs int DryRun bool - // Exec runs skmanage and the socket query. Nil means locally; the kube - // variant supplies a pod-exec so both run inside the router container. + // Exec runs skmanage and the socket query. Exec Execer // SkmanageExtraArgs is appended to every skmanage invocation — e.g. // --ssl-certificate/--ssl-key/--ssl-trustfile when the management @@ -32,7 +31,7 @@ type Result struct { Failed int } -// Run ties the two independent stages together: Gather (gather.go) collects +// Run ties the stages together: Gather (gather.go) collects // raw router + kernel state, Evaluate (criteria.go) applies the idle-time // criteria against that state, and killAll (kill.go) carries out whatever // Evaluate decided. From 5a88b1474aec19c368d2c55528cce59721d89c72 Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Tue, 21 Jul 2026 10:23:07 -0400 Subject: [PATCH 6/9] Changed dry-run to be execute, now must specify to kill connections --- internal/cmd/skupper/common/flags.go | 2 +- internal/cmd/skupper/debug/debug.go | 2 +- internal/cmd/skupper/debug/kube/conn_sweeper.go | 2 +- internal/cmd/skupper/debug/nonkube/conn_sweeper.go | 2 +- internal/cmd/skupper/debug/sweeper/sweeper.go | 6 +++--- 5 files changed, 7 insertions(+), 7 deletions(-) diff --git a/internal/cmd/skupper/common/flags.go b/internal/cmd/skupper/common/flags.go index 5681e56d4..f833e79e6 100644 --- a/internal/cmd/skupper/common/flags.go +++ b/internal/cmd/skupper/common/flags.go @@ -254,7 +254,7 @@ type CommandDebugFlags struct { type CommandConnSweeperFlags struct { IdleThreshold int - DryRun bool + Execute bool } type CommandSystemUninstallFlags struct { diff --git a/internal/cmd/skupper/debug/debug.go b/internal/cmd/skupper/debug/debug.go index 2dc90c8db..07a4736f7 100644 --- a/internal/cmd/skupper/debug/debug.go +++ b/internal/cmd/skupper/debug/debug.go @@ -43,7 +43,7 @@ via adminStatus=deleted.`, var cmdFlags common.CommandConnSweeperFlags cmd.Flags().IntVar(&cmdFlags.IdleThreshold, "idle-threshold", sweeper.DefaultIdleThreshold, "Seconds with no data received before a connection is flagged as orphaned") - cmd.Flags().BoolVar(&cmdFlags.DryRun, "dry-run", false, "List idle connections without killing them") + cmd.Flags().BoolVar(&cmdFlags.Execute, "execute", false, "Close the idle connections found; without this flag they are only listed") kubeCommand.CobraCmd = cmd kubeCommand.Flags = &cmdFlags diff --git a/internal/cmd/skupper/debug/kube/conn_sweeper.go b/internal/cmd/skupper/debug/kube/conn_sweeper.go index 2561010d0..86c494c7a 100644 --- a/internal/cmd/skupper/debug/kube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/kube/conn_sweeper.go @@ -111,7 +111,7 @@ func (cmd *CmdConnSweeper) Run() error { URL: sweeper.DefaultURL, Skmanage: sweeper.DefaultSkmanage, IdleThresholdSecs: cmd.Flags.IdleThreshold, - DryRun: cmd.Flags.DryRun, + Execute: cmd.Flags.Execute, Exec: cmd.podExecer(podName), }) if err != nil { diff --git a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go index 2cd9a35d5..111c27848 100644 --- a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go @@ -92,7 +92,7 @@ func (cmd *CmdConnSweeper) Run() error { URL: cmd.url, Skmanage: cmd.skmanage, IdleThresholdSecs: cmd.Flags.IdleThreshold, - DryRun: cmd.Flags.DryRun, + Execute: cmd.Flags.Execute, Exec: cmd.exec, SkmanageExtraArgs: cmd.sslArgs, }) diff --git a/internal/cmd/skupper/debug/sweeper/sweeper.go b/internal/cmd/skupper/debug/sweeper/sweeper.go index 824407dbc..fa8a83827 100644 --- a/internal/cmd/skupper/debug/sweeper/sweeper.go +++ b/internal/cmd/skupper/debug/sweeper/sweeper.go @@ -15,7 +15,7 @@ type Config struct { URL string Skmanage string IdleThresholdSecs int - DryRun bool + Execute bool // Exec runs skmanage and the socket query. Exec Execer // SkmanageExtraArgs is appended to every skmanage invocation — e.g. @@ -52,8 +52,8 @@ func Run(cfg Config) (Result, error) { return Result{Total: len(snap.TCPConns)}, nil } - if cfg.DryRun { - logf("DRY RUN — would kill %d connection(s):", len(toKill)) + if !cfg.Execute { + logf("Found %d idle connection(s) — re-run with --execute to close them:", len(toKill)) for _, d := range toKill { fmt.Printf(" id=%-6s host=%-25s dir=%s uptime=%-10s reason=%s\n", d.Conn.Identity, d.Conn.Host, d.Conn.Dir, fmtSeconds(d.Conn.UptimeSeconds), d.Reason) From 0fe5de365d87342268c725e10cf86a294d6ed023 Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Tue, 21 Jul 2026 11:06:29 -0400 Subject: [PATCH 7/9] Propagate connection-deletion failures to the CLI + validate idle threshold for upperbound time limit --- internal/cmd/skupper/debug/kube/conn_sweeper.go | 6 ++++++ internal/cmd/skupper/debug/nonkube/conn_sweeper.go | 13 +++++++++++-- internal/cmd/skupper/debug/sweeper/sweeper.go | 2 ++ 3 files changed, 19 insertions(+), 2 deletions(-) diff --git a/internal/cmd/skupper/debug/kube/conn_sweeper.go b/internal/cmd/skupper/debug/kube/conn_sweeper.go index 86c494c7a..9f398a743 100644 --- a/internal/cmd/skupper/debug/kube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/kube/conn_sweeper.go @@ -84,6 +84,9 @@ func (cmd *CmdConnSweeper) ValidateInput(args []string) error { if ok, err := numberValidator.Evaluate(cmd.Flags.IdleThreshold); !ok { validationErrors = append(validationErrors, fmt.Errorf("idle-threshold is not valid: %s", err)) } + if int64(cmd.Flags.IdleThreshold) > sweeper.MaxIdleThreshold { + validationErrors = append(validationErrors, fmt.Errorf("idle-threshold is too large (max %d seconds)", sweeper.MaxIdleThreshold)) + } return errors.Join(validationErrors...) } @@ -132,6 +135,9 @@ func (cmd *CmdConnSweeper) Run() error { if len(failedPods) > 0 { return fmt.Errorf("sweep failed on %d of %d router pod(s): %v", len(failedPods), len(podNames), failedPods) } + if total.Failed > 0 { + return fmt.Errorf("%d idle connection(s) failed to close (%d closed)", total.Failed, total.Killed) + } return nil } diff --git a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go index 111c27848..1dddda850 100644 --- a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go @@ -48,6 +48,9 @@ func (cmd *CmdConnSweeper) ValidateInput(args []string) error { if ok, err := numberValidator.Evaluate(cmd.Flags.IdleThreshold); !ok { validationErrors = append(validationErrors, fmt.Errorf("idle-threshold is not valid: %s", err)) } + if int64(cmd.Flags.IdleThreshold) > sweeper.MaxIdleThreshold { + validationErrors = append(validationErrors, fmt.Errorf("idle-threshold is too large (max %d seconds)", sweeper.MaxIdleThreshold)) + } return errors.Join(validationErrors...) } @@ -88,7 +91,7 @@ func (cmd *CmdConnSweeper) Run() error { if cmd.exec != nil { fmt.Printf("running against %s container %s-skupper-router\n", cmd.platform, cmd.namespace) } - _, err := sweeper.Run(sweeper.Config{ + res, err := sweeper.Run(sweeper.Config{ URL: cmd.url, Skmanage: cmd.skmanage, IdleThresholdSecs: cmd.Flags.IdleThreshold, @@ -96,7 +99,13 @@ func (cmd *CmdConnSweeper) Run() error { Exec: cmd.exec, SkmanageExtraArgs: cmd.sslArgs, }) - return err + if err != nil { + return err + } + if res.Failed > 0 { + return fmt.Errorf("%d idle connection(s) failed to close (%d closed)", res.Failed, res.Killed) + } + return nil } func (cmd *CmdConnSweeper) WaitUntil() error { return nil } diff --git a/internal/cmd/skupper/debug/sweeper/sweeper.go b/internal/cmd/skupper/debug/sweeper/sweeper.go index fa8a83827..d00332506 100644 --- a/internal/cmd/skupper/debug/sweeper/sweeper.go +++ b/internal/cmd/skupper/debug/sweeper/sweeper.go @@ -2,6 +2,7 @@ package sweeper import ( "fmt" + "math" "time" ) @@ -9,6 +10,7 @@ const ( DefaultURL = "amqp://127.0.0.1:5672" DefaultSkmanage = "skmanage" DefaultIdleThreshold = 4 * 3600 // 4 hours + MaxIdleThreshold = math.MaxInt64 / int64(time.Second) ) type Config struct { From 715b7f646829daff35c073c29d983e55a3af27ea Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Wed, 22 Jul 2026 11:43:24 -0400 Subject: [PATCH 8/9] Report socket-state failures instead of reporting zero idle connections --- internal/cmd/skupper/debug/sweeper/gather.go | 23 +++++++++++++------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/internal/cmd/skupper/debug/sweeper/gather.go b/internal/cmd/skupper/debug/sweeper/gather.go index 9fe9e77f6..c661b4f18 100644 --- a/internal/cmd/skupper/debug/sweeper/gather.go +++ b/internal/cmd/skupper/debug/sweeper/gather.go @@ -82,7 +82,10 @@ func Gather(execFn Execer, skmanageBin, url string, extraArgs ...string) (Snapsh } } - byPeer, byLocal := gatherSockets(execFn) + byPeer, byLocal, err := gatherSockets(execFn) + if err != nil { + return Snapshot{}, err + } return Snapshot{ Now: time.Now(), TCPConns: tcpConns, @@ -98,15 +101,19 @@ func isTCPAdaptorConn(c connInfo) bool { // gatherSockets reads kernel socket state, preferring `ss -tin` and falling // back to the python netlink script when ss isn't available (e.g. inside the // router container, which ships python3 but not iproute). -func gatherSockets(execFn Execer) (byPeer, byLocal map[string]socketInfo) { - if out, err := execFn([]string{"ss", "-tin"}); err == nil { - return socketsFromSS(out) +func gatherSockets(execFn Execer) (map[string]socketInfo, map[string]socketInfo, error) { + out, ssErr := execFn([]string{"ss", "-tin"}) + if ssErr == nil { + byPeer, byLocal := socketsFromSS(out) + return byPeer, byLocal, nil } - out, err := execFn([]string{"python3", "-c", inetDiagScript}) - if err != nil { - return map[string]socketInfo{}, map[string]socketInfo{} + + out, pyErr := execFn([]string{"python3", "-c", inetDiagScript}) + if pyErr != nil { + return nil, nil, fmt.Errorf("could not read socket state, so no connection could be matched to its socket: ss unavailable (%v) and python3 fallback failed (%v)", ssErr, pyErr) } - return socketsFromDiagOutput(out) + byPeer, byLocal := socketsFromDiagOutput(out) + return byPeer, byLocal, nil } // socketsFromSS builds two {lastrcv, lastsnd} maps — one keyed by peer From ad001d74575a25aee5fdad859fa7b605dc6f9edd Mon Sep 17 00:00:00 2001 From: Evan Wang Date: Thu, 23 Jul 2026 15:23:26 -0400 Subject: [PATCH 9/9] Fixed some of code rabbits' comments --- .../cmd/skupper/debug/nonkube/conn_sweeper.go | 6 +++++ .../cmd/skupper/debug/sweeper/inetdiag.go | 12 +++++++-- internal/cmd/skupper/debug/sweeper/kill.go | 26 ++++++++++++------- internal/cmd/skupper/debug/sweeper/sweeper.go | 3 +++ 4 files changed, 36 insertions(+), 11 deletions(-) diff --git a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go index 1dddda850..7ded5d498 100644 --- a/internal/cmd/skupper/debug/nonkube/conn_sweeper.go +++ b/internal/cmd/skupper/debug/nonkube/conn_sweeper.go @@ -30,6 +30,7 @@ type CmdConnSweeper struct { skmanage string sslArgs []string exec sweeper.Execer + setupErr error } func NewCmdConnSweeper() *CmdConnSweeper { @@ -62,12 +63,14 @@ func (cmd *CmdConnSweeper) InputToOptions() { platformLoader := &nonkubecommon.NamespacePlatformLoader{} platform, err := platformLoader.Load(cmd.namespace) if err != nil { + cmd.setupErr = fmt.Errorf("could not load platform for namespace %q: %w", cmd.namespace, err) return } cmd.platform = platform url, err := runtime.GetLocalRouterAddress(cmd.namespace) if err != nil { + cmd.setupErr = fmt.Errorf("could not determine router address for namespace %q: %w", cmd.namespace, err) return } cmd.url = url @@ -85,6 +88,9 @@ func (cmd *CmdConnSweeper) InputToOptions() { } func (cmd *CmdConnSweeper) Run() error { + if cmd.setupErr != nil { + return cmd.setupErr + } if cmd.url == "" || cmd.skmanage == "" { return fmt.Errorf("could not determine router management address for namespace %q", cmd.namespace) } diff --git a/internal/cmd/skupper/debug/sweeper/inetdiag.go b/internal/cmd/skupper/debug/sweeper/inetdiag.go index 4498deaba..96ba58a2c 100644 --- a/internal/cmd/skupper/debug/sweeper/inetdiag.go +++ b/internal/cmd/skupper/debug/sweeper/inetdiag.go @@ -16,7 +16,7 @@ import ( // TODO: IPv4 only; extend for IPv6. const inetDiagScript = ` -import socket, struct +import socket, struct, sys, os NETLINK_SOCK_DIAG = 4 SOCK_DIAG_BY_FAMILY = 20 @@ -39,7 +39,15 @@ while not done: off = 0 while off + 16 <= len(data): ln, typ = struct.unpack_from("=IH", data, off) - if ln < 16 or typ in (NLMSG_DONE, NLMSG_ERROR): + if ln < 16 or typ == NLMSG_DONE: + done = True + break + if typ == NLMSG_ERROR: + # nlmsgerr starts with a signed errno (negative on failure, 0 = ACK) + err = struct.unpack_from("=i", data, off + 16)[0] + if err != 0: + sys.stderr.write("netlink SOCK_DIAG error: %s\n" % os.strerror(-err)) + sys.exit(1) done = True break # inet_diag_msg: sport/dport are big-endian at +4/+6, src at +8, diff --git a/internal/cmd/skupper/debug/sweeper/kill.go b/internal/cmd/skupper/debug/sweeper/kill.go index 21d15be58..9838e7c05 100644 --- a/internal/cmd/skupper/debug/sweeper/kill.go +++ b/internal/cmd/skupper/debug/sweeper/kill.go @@ -1,5 +1,7 @@ package sweeper +import "strings" + // killAll force-closes each flagged connection by setting adminStatus=deleted // via skmanage. func killAll(execFn Execer, skmanageBin, url string, extraArgs []string, decisions []Decision) (killed, failed int) { @@ -18,19 +20,25 @@ func killAll(execFn Execer, skmanageBin, url string, extraArgs []string, decisio continue } // Killing one half of a proxied pair cascade-closes the other half, - // so a later kill of that half fails with "not found". Confirm the - // connection is really gone before calling it a failure. - if _, readErr := runSkmanage(execFn, skmanageBin, url, extraArgs, - "READ", "--type="+ConnType, "--identity="+d.Conn.Identity); readErr != nil { + // so a later kill of that half fails with "not found". + _, readErr := runSkmanage(execFn, skmanageBin, url, extraArgs, + "READ", "--type="+ConnType, "--identity="+d.Conn.Identity) + if readErr != nil && isNotFound(readErr) { logf(" id=%s host=%s dir=%s reason=%s → already closed", d.Conn.Identity, d.Conn.Host, d.Conn.Dir, d.Reason) killed++ - } else { - logf(" id=%s host=%s dir=%s reason=%s → failed: %s", - d.Conn.Identity, d.Conn.Host, d.Conn.Dir, - d.Reason, err.Error()) - failed++ + continue } + logf(" id=%s host=%s dir=%s reason=%s → failed: %s", + d.Conn.Identity, d.Conn.Host, d.Conn.Dir, d.Reason, err.Error()) + failed++ } return killed, failed } + +// isNotFound reports whether a skmanage READ error indicates the connection no +// longer exists rather than the router being unreachable or some other +// management failure that leaves the connection's fate unknown. +func isNotFound(err error) bool { + return strings.Contains(strings.ToLower(err.Error()), "not found") +} diff --git a/internal/cmd/skupper/debug/sweeper/sweeper.go b/internal/cmd/skupper/debug/sweeper/sweeper.go index d00332506..d61a511bc 100644 --- a/internal/cmd/skupper/debug/sweeper/sweeper.go +++ b/internal/cmd/skupper/debug/sweeper/sweeper.go @@ -38,6 +38,9 @@ type Result struct { // criteria against that state, and killAll (kill.go) carries out whatever // Evaluate decided. func Run(cfg Config) (Result, error) { + if cfg.IdleThresholdSecs < 0 || int64(cfg.IdleThresholdSecs) > MaxIdleThreshold { + return Result{}, fmt.Errorf("idle threshold must be between 0 and %d seconds", MaxIdleThreshold) + } if cfg.Exec == nil { cfg.Exec = LocalExec }