diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d68738e..019d5dd 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -34,11 +34,11 @@ jobs: # through a tailcat destination dials via tailcat.Client, the same # way tailcat's own "forward"/"socks" subcommands do), so go vet/go # test need that directory to exist before either can even resolve - # imports. A plain, unpatched clone is enough here: the two patches - # below only change Android-specific networking behavior, irrelevant - # to vetting/testing meowshell on this runner's own host platform. - # ./build.sh (further down) later reuses and patches this same - # checkout for the real cross-compiled build. + # imports. Apply the reviewed live-path API patch before compiling + # meowshell: the agent now uses that exact Tailcat Client directly, so + # path reporting is statically type-checked rather than hidden behind + # reflection. ./build.sh later resets this disposable checkout and + # reapplies the full reviewed patch set before release binaries are built. - name: Fetch tailcat source (for go.mod's replace directive) env: SRC_REF: ${{ inputs.tailcat_ref }} @@ -47,6 +47,7 @@ jobs: git clone --depth=1 https://github.com/tailscale/tailcat.git .tailcat-src git -C .tailcat-src fetch --depth=1 origin "$SRC_REF" git -C .tailcat-src checkout --detach FETCH_HEAD + git -C .tailcat-src apply "$GITHUB_WORKSPACE/patches/tailcat/live-path-status.patch" - name: Vet and test meowshell env: @@ -210,7 +211,7 @@ jobs: steps: - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 - - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 + - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v5 with: go-version: '1.27.1' @@ -242,7 +243,7 @@ jobs: steps: - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 - - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 + - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v5 with: go-version: '1.27.1' @@ -258,6 +259,7 @@ jobs: git clone --depth=1 https://github.com/tailscale/tailcat.git .tailcat-src git -C .tailcat-src fetch --depth=1 origin "$SRC_REF" git -C .tailcat-src checkout --detach FETCH_HEAD + git -C .tailcat-src apply "$GITHUB_WORKSPACE/patches/tailcat/live-path-status.patch" # Some of cmd/meowshell's Windows-specific code (known_hosts DACL/SID # validation, N9) only builds and runs under GOOS=windows, so it never @@ -332,7 +334,7 @@ jobs: with: dotnet-version: "8.0.x" - - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 + - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v5 with: go-version: '1.27.1' @@ -356,6 +358,9 @@ jobs: dotnet restore dotnet/Meowshell.sln --force --no-cache --nologo dotnet restore scripts/CommentStripper/CommentStripper.csproj --force --no-cache --nologo + - name: Start a shared local DERP relay for direct-path E2E + run: ./e2e/start-testderp.sh + - name: Test # The E2E tests in Meowshell.Tests run against these real binaries # instead of the fake stand-ins the rest of the suite uses; they @@ -426,6 +431,10 @@ jobs: path: nupkg/* if-no-files-found: error + - name: Stop the shared local DERP relay + if: always() + run: kill "$(cat /tmp/testderp.pid)" 2>/dev/null || true + # Meowshell.Demo: a real, installable app, not just a pass/fail check. # One button generates a fresh throwaway shell address, one field copies # it. Published as a single universal APK bundling all four ABIs, the diff --git a/cmd/meowshell/agent.go b/cmd/meowshell/agent.go index 802e0c4..001332b 100644 --- a/cmd/meowshell/agent.go +++ b/cmd/meowshell/agent.go @@ -7,6 +7,7 @@ import ( "flag" "fmt" "io" + "log" "net" "os" "path/filepath" @@ -20,6 +21,7 @@ import ( "golang.org/x/crypto/ssh" "golang.org/x/crypto/ssh/agent" "tailscale.com/types/key" + "tailscale.com/types/logger" ) const agentUsage = `meowshell agent -- a persistent, multiplexed SSH connection @@ -33,8 +35,9 @@ protocol.go) instead of the one-process-per-operation model "meowshell connect"/"cp" use. Opening a shell and running a command against the same host costs one login instead of two. - is a tailcat address (dialed through tailcat's own bare -client mode, same as "connect"/"cp") or a "[user@]host[:port]" TCP address + is a tailcat address (dialed by the in-process tailcat client +so the same live WireGuard path can be observed and reused for forwarding) +or a "[user@]host[:port]" TCP address for a general (non-tailcat) SSH host, verified against a known_hosts file (--known-hosts) with trust-on-first-use for a host seen for the first time. --jump chains through one or more intermediate TCP hosts first, each @@ -162,6 +165,14 @@ func agentCmd(args []string) error { if err := session.writeControl(0, connected); err != nil { return err } + + pathCtx, cancelPath := context.WithCancel(context.Background()) + defer cancelPath() + if source := session.tailcatPathSource(); source != nil { + go reportTailcatPath(pathCtx, source, agentPathPollInterval, func(msg controlMessage) error { + return session.writeControl(0, msg) + }) + } return <-frameErrCh } @@ -351,16 +362,32 @@ func (a *agentSession) connect(ctx context.Context, opts connectOptions) error { user := "" if last && looksLikeTailcatAddress(hop) { - dial = tailcatDialer(opts.tailcatBin, tailcatClientArgv(opts.key, opts.derpMapURL, opts.verbose, hop, opts.port)) hkCallback = tailcatHostKeyCallback() tcKey, err := tailcatKeyFromName(opts.key) if err != nil { closeClients(chain) - return fmt.Errorf("resolving --key %q for forwarding: %w", opts.key, err) + return fmt.Errorf("resolving --key %q for Tailcat: %w", opts.key, err) + } + tcClient := &tailcat.Client{ + Server: tailcat.Addr(hop), + Key: tcKey, + DERPMapURL: opts.derpMapURL, + Logf: logger.Discard, + } + if opts.verbose { + tcClient.Logf = log.Printf + } + dial, err = tailcatClientDialer(tcClient, opts.port) + if err != nil { + closeClients(chain) + return err } a.tcAddr = tailcat.Addr(hop) a.tcKey = tcKey a.tcDERPMapURL = opts.derpMapURL + a.tcMu.Lock() + a.tcClient = tcClient + a.tcMu.Unlock() } else { var hostPort string user, hostPort = splitUserHost(hop, opts.port) @@ -389,6 +416,9 @@ func (a *agentSession) connect(ctx context.Context, opts connectOptions) error { sc, err := dialSSHClient(ctx, dial, remoteAddr, user, hkCallback, opts.auth) if err != nil { closeClients(chain) + if last && looksLikeTailcatAddress(hop) { + a.closeTailcatClient() + } return fmt.Errorf("connecting to %s: %w", connectTargetForDiagnostics(hop, last && looksLikeTailcatAddress(hop)), err) } chain = append(chain, sc) @@ -421,12 +451,25 @@ func (a *agentSession) closeHops() { // after which closing/clearing the client is bounded. closeClients(a.hops) a.resetSFTPClient(nil) + a.closeTailcatClient() +} + +func (a *agentSession) closeTailcatClient() { a.tcMu.Lock() + defer a.tcMu.Unlock() if a.tcClient != nil { a.tcClient.Close() a.tcClient = nil } - a.tcMu.Unlock() +} + +func (a *agentSession) tailcatPathSource() *tailcatForwardClient { + a.tcMu.Lock() + defer a.tcMu.Unlock() + if a.tcClient == nil { + return nil + } + return &tailcatForwardClient{cl: a.tcClient} } const ( diff --git a/cmd/meowshell/agent_e2e_test.go b/cmd/meowshell/agent_e2e_test.go index c879eeb..9cb7560 100644 --- a/cmd/meowshell/agent_e2e_test.go +++ b/cmd/meowshell/agent_e2e_test.go @@ -175,7 +175,7 @@ func expectConnected(t *testing.T, r *bufio.Reader) { } } -func expectChannelOpened(t *testing.T, r *bufio.Reader) uint32 { +func expectChannelOpenedMessage(t *testing.T, r *bufio.Reader) (frame, controlMessage) { t.Helper() for { f, err := readFrameWithDeadline(t, r) @@ -191,10 +191,46 @@ func expectChannelOpened(t *testing.T, r *bufio.Reader) uint32 { } switch msg.Msg { case "channel_opened": - return f.ChannelID + return f, msg case "error": t.Fatalf("agent returned an error opening the channel: %s: %s", msg.Code, msg.Message) + case "path": + // Live path telemetry is an asynchronous connection-level + // notification. It may legally arrive between a request and that + // request's reply, so request/response E2E helpers must not consume + // it as the synchronous response they are waiting for. + continue + } + } +} + +func expectChannelOpened(t *testing.T, r *bufio.Reader) uint32 { + t.Helper() + f, _ := expectChannelOpenedMessage(t, r) + return f.ChannelID +} + +func expectRequestReply(t *testing.T, r *bufio.Reader, requestID string) controlMessage { + t.Helper() + for { + f, err := readFrameWithDeadline(t, r) + if err != nil { + t.Fatalf("reading reply for request %q: %v", requestID, err) + } + if f.Type != frameTypeControl { + continue + } + var msg controlMessage + if err := json.Unmarshal(f.Payload, &msg); err != nil { + t.Fatalf("decoding control message: %v", err) + } + if msg.Msg == "path" { + continue + } + if msg.RequestID != requestID { + t.Fatalf("reply RequestID = %q, want %q (msg=%+v)", msg.RequestID, requestID, msg) } + return msg } } diff --git a/cmd/meowshell/agent_path.go b/cmd/meowshell/agent_path.go index 8113ca8..1f1bc50 100644 --- a/cmd/meowshell/agent_path.go +++ b/cmd/meowshell/agent_path.go @@ -19,6 +19,10 @@ type agentPathStatusSource interface { PathStatus() (agentPathStatus, bool) } +type agentPathProber interface { + ProbePath(context.Context) (agentPathStatus, bool) +} + type agentPathState struct { direct bool via string @@ -54,11 +58,7 @@ func reportTailcatPath( var last agentPathState haveLast := false - poll := func() bool { - status, ok := source.PathStatus() - if !ok { - return true - } + emitStatus := func(status agentPathStatus) bool { msg := pathControlMessageFromStatus(status) state := agentPathState{direct: *msg.Direct, via: msg.Via} if haveLast && state == last { @@ -71,6 +71,32 @@ func reportTailcatPath( haveLast = true return true } + poll := func() bool { + status, ok := source.PathStatus() + if ok { + if !emitStatus(status) { + return false + } + if status.direct { + return true + } + } + + // Tailcat's DiscoPing actively nudges the call-me-maybe endpoint + // exchange. Probe while the path is unknown as well as while it is + // relayed: on a freshly started client Status can lag behind the live + // magicsock route, and returning early here used to prevent the very + // probe needed to establish a direct endpoint. + if prober, hasProber := source.(agentPathProber); hasProber { + probeCtx, cancel := context.WithTimeout(ctx, 2*time.Second) + probed, probeOK := prober.ProbePath(probeCtx) + cancel() + if probeOK { + return emitStatus(probed) + } + } + return true + } if !poll() { return diff --git a/cmd/meowshell/agent_path_test.go b/cmd/meowshell/agent_path_test.go index 3a3c1d0..e369899 100644 --- a/cmd/meowshell/agent_path_test.go +++ b/cmd/meowshell/agent_path_test.go @@ -97,3 +97,97 @@ func TestReportTailcatPathIgnoresUnknownAndEmitsOnlyPathChanges(t *testing.T) { t.Fatalf("second path = direct %#v via %q, want direct with no relay", got[1].Direct, got[1].Via) } } + +type probingPathSource struct { + mu sync.Mutex + direct bool + probes int +} + +func (s *probingPathSource) PathStatus() (agentPathStatus, bool) { + s.mu.Lock() + defer s.mu.Unlock() + return agentPathStatus{ + direct: s.direct, + relayRegion: "ci", + }, true +} + +func (s *probingPathSource) ProbePath(context.Context) (agentPathStatus, bool) { + s.mu.Lock() + defer s.mu.Unlock() + s.probes++ + s.direct = true + return agentPathStatus{direct: true, endpoint: "203.0.113.7:41641"}, true +} + +func TestReportTailcatPathProbesRelayedConnectionAndEmitsDirectUpgrade(t *testing.T) { + source := &probingPathSource{} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + var got []controlMessage + reportTailcatPath(ctx, source, time.Millisecond, func(msg controlMessage) error { + got = append(got, msg) + if len(got) == 2 { + cancel() + } + return nil + }) + + if len(got) != 2 { + t.Fatalf("got %d path updates, want 2", len(got)) + } + if got[0].Direct == nil || *got[0].Direct || got[0].Via != "ci" { + t.Fatalf("first update = %#v, want relayed via ci", got[0]) + } + if got[1].Direct == nil || !*got[1].Direct || got[1].Via != "" { + t.Fatalf("second update = %#v, want direct", got[1]) + } + source.mu.Lock() + defer source.mu.Unlock() + if source.probes == 0 { + t.Fatal("relayed path was never actively probed") + } +} + +type initiallyUnknownProbingPathSource struct { + mu sync.Mutex + probes int +} + +func (s *initiallyUnknownProbingPathSource) PathStatus() (agentPathStatus, bool) { + return agentPathStatus{}, false +} + +func (s *initiallyUnknownProbingPathSource) ProbePath(context.Context) (agentPathStatus, bool) { + s.mu.Lock() + defer s.mu.Unlock() + s.probes++ + return agentPathStatus{direct: true, endpoint: "127.0.0.1:41641"}, true +} + +func TestReportTailcatPathProbesUnknownConnectionAndUsesLivePingResult(t *testing.T) { + source := &initiallyUnknownProbingPathSource{} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + var got []controlMessage + reportTailcatPath(ctx, source, time.Millisecond, func(msg controlMessage) error { + got = append(got, msg) + cancel() + return nil + }) + + if len(got) != 1 { + t.Fatalf("got %d path updates, want 1", len(got)) + } + if got[0].Direct == nil || !*got[0].Direct || got[0].Via != "" { + t.Fatalf("path update = %#v, want direct probe result", got[0]) + } + source.mu.Lock() + defer source.mu.Unlock() + if source.probes == 0 { + t.Fatal("unknown path was never actively probed") + } +} diff --git a/cmd/meowshell/agent_sftp_e2e_test.go b/cmd/meowshell/agent_sftp_e2e_test.go index 9406c50..95e28ed 100644 --- a/cmd/meowshell/agent_sftp_e2e_test.go +++ b/cmd/meowshell/agent_sftp_e2e_test.go @@ -114,8 +114,7 @@ func TestAgentSFTPEndToEnd(t *testing.T) { uploadViaAgent(t, stdin, out, "big.bin", payload, false, 0, 0) send(t, stdin, 0, controlMessage{Msg: "open_channel", Kind: "sftp_download", Path: "big.bin"}) - f := mustReadFrame(t, out) - opened := decodeControl(t, f) + f, opened := expectChannelOpenedMessage(t, out) if opened.Msg != "channel_opened" || opened.Size != int64(len(payload)) { t.Fatalf("channel_opened = %+v, want Size %d", opened, len(payload)) } @@ -193,11 +192,7 @@ func sftpOp(t *testing.T, stdin interface { req.Msg = "sftp_op" req.RequestID = fmt.Sprintf("op%d", time.Now().UnixNano()) send(t, stdin, 0, req) - f := mustReadFrame(t, out) - msg := decodeControl(t, f) - if msg.RequestID != req.RequestID { - t.Fatalf("sftp_op %s: reply RequestID = %q, want %q (msg=%+v)", req.Op, msg.RequestID, req.RequestID, msg) - } + msg := expectRequestReply(t, out, req.RequestID) if msg.Msg == "error" { t.Fatalf("sftp_op %s %s failed: %s: %s", req.Op, req.Path, msg.Code, msg.Message) } @@ -211,8 +206,7 @@ func trySFTPOp(t *testing.T, stdin interface { req.Msg = "sftp_op" req.RequestID = fmt.Sprintf("op%d", time.Now().UnixNano()) send(t, stdin, 0, req) - f := mustReadFrame(t, out) - msg := decodeControl(t, f) + msg := expectRequestReply(t, out, req.RequestID) if msg.Msg == "error" { return msg.Code } diff --git a/cmd/meowshell/agent_tailcat_forward_e2e_test.go b/cmd/meowshell/agent_tailcat_forward_e2e_test.go index 58aab63..1f3ef17 100644 --- a/cmd/meowshell/agent_tailcat_forward_e2e_test.go +++ b/cmd/meowshell/agent_tailcat_forward_e2e_test.go @@ -77,8 +77,7 @@ func TestAgentForwardsThroughTailcatDestination(t *testing.T) { Msg: "open_channel", Kind: "forward_local", ListenAddr: "127.0.0.1:0", RemoteAddr: "localhost:" + backendPort, }) - f := mustReadFrame(t, out) - opened := decodeControl(t, f) + _, opened := expectChannelOpenedMessage(t, out) if opened.Msg != "channel_opened" || opened.BoundAddr == "" { t.Fatalf("channel_opened = %+v, want a non-empty BoundAddr", opened) } @@ -100,8 +99,14 @@ func TestAgentForwardsThroughTailcatDestination(t *testing.T) { t.Run("forward_socks", func(t *testing.T) { send(t, stdin, 0, controlMessage{Msg: "open_channel", Kind: "forward_socks", ListenAddr: "127.0.0.1:0"}) - f := mustReadFrame(t, out) - opened := decodeControl(t, f) + // Use the path-telemetry-aware helper, not a raw mustReadFrame: this + // server has a live path reporter goroutine (TS_DEBUG_TAILCAT_LOCAL_DERP + // above), which can interleave an asynchronous "path" control message + // between this request and its channel_opened reply. forward_local + // already goes through expectChannelOpenedMessage for the same reason; + // this subtest used to read the raw next frame instead and flaked + // whenever a path notification won the race. + _, opened := expectChannelOpenedMessage(t, out) if opened.Msg != "channel_opened" || opened.BoundAddr == "" { t.Fatalf("channel_opened = %+v, want a non-empty BoundAddr", opened) } diff --git a/cmd/meowshell/netmon_android.go b/cmd/meowshell/netmon_android.go new file mode 100644 index 0000000..d73aad3 --- /dev/null +++ b/cmd/meowshell/netmon_android.go @@ -0,0 +1,57 @@ +// Copyright (c) Tailscale Inc & contributors +// SPDX-License-Identifier: BSD-3-Clause + +//go:build android + +package main + +import ( + "github.com/wlynxg/anet" + "tailscale.com/net/netmon" +) + +// Go's net.Interfaces() and net.Interface.Addrs() open a netlink socket and +// call bind() on it explicitly; Android's SELinux policy denies that bind +// operation on a netlink_route_socket to every app in the untrusted_app +// domain. tailcat.Client.initLocked calls netmon.New() unconditionally +// before it can dial anything, so with no working interface getter that +// call fails outright with "netmon.New: route ip+net: netlinkrib: +// permission denied" -- the exact error PROBE_SSH_FAIL reported once this +// repo started dialing Tailcat addresses this way. +// +// Before the in-process tailcat.Client added in this branch, every Tailcat +// SSH connection went through the separate tailcat subprocess, and +// cmd/tailcat/netmon_android.go (added by +// patches/tailcat/android-netmon-interface-getter.patch) already registered +// this same workaround there -- so netmon.New() never ran unpatched inside +// an Android process. Forwarding's tailcatClientDialer and agent.connect's +// direct-address path now call netmon.New() in *this* binary instead, and +// that patch only touches the tailcat module, not meowshell's, so the +// meowshell binary needs its own copy of the registration or every direct +// Tailcat connection on Android regresses back to the permission denial. +// +// anet reimplements the same netlink RIB queries the stdlib does, but never +// binds the socket explicitly -- it lets the kernel's implicit +// autobind-on-send handle that instead, which Android's policy does not +// deny -- so it returns real, non-empty interface data where the stdlib +// call fails outright. +func init() { + netmon.RegisterInterfaceGetter(func() ([]netmon.Interface, error) { + ifs, err := anet.Interfaces() + if err != nil { + return nil, err + } + ret := make([]netmon.Interface, len(ifs)) + for i := range ifs { + // AltAddrs, not a second RegisterInterfaceGetter-style hook: + // netmon.Interface.Addrs() only consults i.Interface.Addrs() + // (the same netlink-based stdlib call) when AltAddrs is nil, + // so leaving it unset here would hit the identical permission + // denial one level down, per interface, right after fixing the + // enumeration itself. + addrs, _ := anet.InterfaceAddrsByInterface(&ifs[i]) + ret[i] = netmon.Interface{Interface: &ifs[i], AltAddrs: addrs} + } + return ret, nil + }) +} diff --git a/cmd/meowshell/tailcatdial.go b/cmd/meowshell/tailcatdial.go index cd3d4dd..ebd013f 100644 --- a/cmd/meowshell/tailcatdial.go +++ b/cmd/meowshell/tailcatdial.go @@ -20,6 +20,48 @@ type tailcatForwardClient struct { cl *tailcat.Client } +func tailcatClientDialer(cl *tailcat.Client, port string) (dialer, error) { + p, err := strconv.ParseUint(port, 10, 16) + if err != nil || p == 0 { + return nil, fmt.Errorf("invalid Tailcat SSH port %q", port) + } + return func(ctx context.Context) (net.Conn, error) { + return cl.DialTCPPort(ctx, uint16(p)) + }, nil +} + +func (c *tailcatForwardClient) PathStatus() (agentPathStatus, bool) { + status, ok := c.cl.PathStatus() + if !ok { + return agentPathStatus{}, false + } + return agentPathStatus{ + direct: status.Direct, + endpoint: status.Endpoint, + relayRegion: status.RelayRegion, + txBytes: status.TxBytes, + rxBytes: status.RxBytes, + }, true +} + +func (c *tailcatForwardClient) ProbePath(ctx context.Context) (agentPathStatus, bool) { + result, err := c.cl.DiscoPing(ctx) + if err != nil || result == nil { + return agentPathStatus{}, false + } + status := agentPathStatus{ + direct: result.Endpoint != "", + endpoint: result.Endpoint, + } + if !status.direct { + status.relayRegion = result.DERPRegionCode + if status.relayRegion == "" { + status.relayRegion = fmt.Sprint(result.DERPRegionID) + } + } + return status, true +} + func (c *tailcatForwardClient) Dial(network, addr string) (net.Conn, error) { if network != "tcp" { return nil, fmt.Errorf("tailcat forwarding only supports tcp, not %q", network) diff --git a/cmd/meowshell/tailcatdial_test.go b/cmd/meowshell/tailcatdial_test.go index c9089d5..3a224bd 100644 --- a/cmd/meowshell/tailcatdial_test.go +++ b/cmd/meowshell/tailcatdial_test.go @@ -88,3 +88,4 @@ func TestTailcatKeyFromNameRejectsOversizedKeyFile(t *testing.T) { t.Fatal("oversized key file was accepted") } } + diff --git a/dotnet/Meowshell.Tests/MeowshellAgentConnectionE2ETests.cs b/dotnet/Meowshell.Tests/MeowshellAgentConnectionE2ETests.cs index 034f1de..d2bef56 100644 --- a/dotnet/Meowshell.Tests/MeowshellAgentConnectionE2ETests.cs +++ b/dotnet/Meowshell.Tests/MeowshellAgentConnectionE2ETests.cs @@ -3,6 +3,7 @@ using System.Net.Sockets; using System.Text; using System.Text.RegularExpressions; +using System.Text.Json; using Meowshell; namespace Meowshell.Tests; @@ -19,6 +20,7 @@ private static void Mask(string value) private const string TailcatEnvVar = "DOTNET_E2E_TAILCAT_BIN"; private const string MeowshellEnvVar = "DOTNET_E2E_MEOWSHELL_BIN"; + private const string RelayMetricsEnvVar = "TESTDERP_METRICS_URL"; private readonly string _dir = Directory.CreateTempSubdirectory("agent-connection-e2e-").FullName; @@ -373,6 +375,83 @@ public async Task SessionLogCallbackSeesTheConnectionSetupsOwnDiagnostics() } } + [Fact] + public async Task DirectPathLargePayloadBypassesTheRelay() + { + var real = FindRealBinaries(); + if (real is null) return; + var (bin, _) = real.Value; + + var metricsUrl = Environment.GetEnvironmentVariable(RelayMetricsEnvVar); + Assert.False(string.IsNullOrWhiteSpace(metricsUrl), + $"Real-binary E2E requires {RelayMetricsEnvVar}; CI must start e2e/testderp before this test."); + + // Deliberately do not use RelayE2E.StartServerAsync here. That helper + // gives each server its own embedded loopback DERP. This test needs + // both endpoints on the shared testderp instance so its counters are + // an independent upper bound on how much payload traversed the relay. + await using var server = await MeowshellServer.StartAsync(new MeowshellOptions + { + BinaryDirectory = bin, + HomeDirectory = Path.Combine(_dir, "direct-server-home"), + WorkDirectory = Path.Combine(_dir, "direct-server-work"), + InsecureNoAuth = true, + Lifetime = TimeSpan.FromMinutes(2), + StartTimeout = TimeSpan.FromSeconds(30), + }); + Mask(server.Address); + + await using var connection = await MeowshellAgentConnection.ConnectAsync(ClientOptions(bin), server.Address); + + var becameDirect = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + void OnPath(MeowshellPathStatus status) + { + if (status.Direct) becameDirect.TrySetResult(status); + } + + connection.PathChanged += OnPath; + try + { + if (connection.CurrentPath?.Direct != true) + await becameDirect.Task.WaitAsync(TimeSpan.FromSeconds(30)); + } + finally + { + connection.PathChanged -= OnPath; + } + + Assert.True(connection.CurrentPath?.Direct == true, + $"Tailcat never established a direct path; last path was {connection.CurrentPath}."); + + // Let any relay-assisted discovery frames settle before taking the + // baseline. Direct data may still coexist with tiny DERP discovery + // traffic; the assertion below deliberately allows a small control + // budget while proving the 1 MiB application payload did not traverse + // the relay. + await Task.Delay(500); + var before = await ReadRelayPayloadBytesAsync(metricsUrl!); + + const int applicationBytes = 1024 * 1024; + await using var exec = await connection.OpenExecAsync( + [$"head -c {applicationBytes} /dev/zero | tr '\\0' 'D'"]); + + long received = 0; + var buffer = new byte[32 * 1024]; + int read; + while ((read = await exec.Output.ReadAsync(buffer)) > 0) + received += read; + + Assert.Equal(0, await exec.Completed); + Assert.Equal(applicationBytes, received); + + var after = await ReadRelayPayloadBytesAsync(metricsUrl!); + var relayDelta = after - before; + + const long maxControlBytes = 64 * 1024; + Assert.InRange(relayDelta, 0, maxControlBytes); + } + [Fact] public async Task ExecChannelDeliversLargeOutputIntactUnderBackpressure() { @@ -406,6 +485,15 @@ public async Task ExecChannelDeliversLargeOutputIntactUnderBackpressure() Assert.All(received, b => Assert.Equal((byte)'A', b)); } + private static async Task ReadRelayPayloadBytesAsync(string metricsUrl) + { + using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(5) }; + using var document = JsonDocument.Parse(await http.GetStringAsync(metricsUrl)); + var root = document.RootElement; + return root.GetProperty("bytes_received").GetInt64() + + root.GetProperty("bytes_sent").GetInt64(); + } + private static async Task Socks5GreetAsync(Socket socket, byte[] methods, CancellationToken cancellationToken) { var greeting = new byte[2 + methods.Length]; diff --git a/e2e/start-testderp.sh b/e2e/start-testderp.sh index 7bbf513..dc960e9 100755 --- a/e2e/start-testderp.sh +++ b/e2e/start-testderp.sh @@ -1,6 +1,7 @@ #!/usr/bin/env bash # Starts e2e/testderp (a single-node, loopback-only DERP relay) in the -# background and exports TAILCAT_DERPMAP_URL for the rest of this job, so +# background and exports TAILCAT_DERPMAP_URL plus TESTDERP_METRICS_URL for +# the rest of this job, so # meowshell/tailcat E2E tests never depend on reaching the public Tailscale # relay infrastructure. Every subsequent step in the job inherits the # variable as a real environment variable (via $GITHUB_ENV), and tailcat @@ -23,7 +24,9 @@ for _ in $(seq 1 100); do # Unset outside Actions: this script is also the documented way to # run an E2E suite locally (see README.md), where there is no # $GITHUB_ENV to export through and set -u would abort here. + metrics_url="${url%/derpmap.json}/metrics.json" echo "TAILCAT_DERPMAP_URL=$url" >> "${GITHUB_ENV:-/dev/null}" + echo "TESTDERP_METRICS_URL=$metrics_url" >> "${GITHUB_ENV:-/dev/null}" echo "local DERP relay ready: $url" exit 0 fi diff --git a/e2e/testderp/main.go b/e2e/testderp/main.go index 0cb8aa7..a98abd6 100644 --- a/e2e/testderp/main.go +++ b/e2e/testderp/main.go @@ -131,6 +131,10 @@ func run(statusFile, verifyClientURL string, verifyClientFailOpen bool) error { w.Header().Set("Content-Type", "application/json") w.Write(mapJSON) }) + mapMux.HandleFunc("/metrics.json", func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = fmt.Fprint(w, d.ExpVar(false).String()) + }) mapSrv := &http.Server{Handler: mapMux} go mapSrv.Serve(mapLn) defer mapSrv.Close() diff --git a/go.mod b/go.mod index 2d66dac..fb2f25c 100644 --- a/go.mod +++ b/go.mod @@ -63,6 +63,7 @@ require ( github.com/tailscale/wireguard-go v0.0.0-20260904023712-e855235c55a2 // indirect github.com/u-root/u-root v0.14.0 // indirect github.com/unixshells/vt-go v0.1.0 // indirect + github.com/wlynxg/anet v0.0.5 // indirect github.com/x448/float16 v0.8.4 // indirect github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect go4.org/mem v0.0.0-20240501181205-ae6ca9944745 // indirect diff --git a/go.sum b/go.sum index 2f44eab..f2cdcb6 100644 --- a/go.sum +++ b/go.sum @@ -215,6 +215,8 @@ github.com/unixshells/vt-go v0.1.0 h1:HWILcExV8MHc6o2YqIRMIM/KrjH31kPzZrlvLS9BG6 github.com/unixshells/vt-go v0.1.0/go.mod h1:dZpOXxVyO7LG1uxkVMixX7HwUf2ZbUkEKhtHHJTTIBI= github.com/vishvananda/netns v0.0.5 h1:DfiHV+j8bA32MFM7bfEunvT8IAqQ/NzSJHtcmW5zdEY= github.com/vishvananda/netns v0.0.5/go.mod h1:SpkAiCQRtJ6TvvxPnOSyH3BMl6unz3xZlaprSwhNNJM= +github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU= +github.com/wlynxg/anet v0.0.5/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no=