diff --git a/docs/architecture.mdx b/docs/architecture.mdx index f7e7ffed..0e669a4c 100644 --- a/docs/architecture.mdx +++ b/docs/architecture.mdx @@ -587,13 +587,17 @@ scanner mints one and persists it. ### Sharing UDP 5353 PAIR runs its own mDNS responder rather than depending on a system one, because -Windows ships none. That responder must coexist with whatever else is on the -port, including Bonjour, Avahi, and PAIR's own sibling processes. It therefore -sets `SO_REUSEADDR` on the socket to share UDP 5353. - -It deliberately does **not** set `SO_REUSEPORT`. On Linux that load-balances -incoming unicast datagrams across every socket sharing the port, which would let -one process swallow mDNS replies meant for another. +Windows ships none. Its receive socket and short-lived per-interface send +sockets bind UDP 5353, as RFC 6762 requires for mDNS queries and responses. +Binding each sender to the selected interface address also preserves reliable +egress on multi-homed Windows hosts. + +Those sockets set `SO_REUSEADDR` so they coexist with Bonjour, Avahi, and other +PAIR processes. A Darwin sender also sets `SO_REUSEPORT`, matching the BSD +multicast sharing behavior needed to coexist with the system mDNS responder. +Linux deliberately omits `SO_REUSEPORT`: there it can load-balance incoming +unicast datagrams into a short-lived send socket and steal replies from the +long-lived receiver. ### Node Enrichment diff --git a/services/readme.md b/services/readme.md index ee3523d5..efc1a9de 100644 --- a/services/readme.md +++ b/services/readme.md @@ -55,7 +55,11 @@ This tree builds thirteen Go binaries. `nvpair-ui-broker` is the parent service Shared code lives in the local `shared/` Go module (imported as `nvpair-shared/…`, replaced via `replace nvpair-shared => ../shared`). It provides logging, wire types, JSON-RPC and IPC, discovery records, mDNS, network monitoring, stable node identity, application data paths, and cluster trust helpers. -The mDNS responder is our own rather than the host's, because Windows ships none. It sets `SO_REUSEADDR` so it shares UDP 5353 with sibling PAIR processes and with a system responder — `avahi-daemon` on Linux, Bonjour where present — needing no configuration on either platform. +The mDNS responder is our own rather than the host's, because Windows ships +none. Its receive and per-interface send sockets bind UDP 5353 as RFC 6762 +requires. Socket reuse lets them coexist with sibling PAIR processes and with a +system responder — `avahi-daemon` on Linux or Bonjour where present — without +configuration. The broker feeds every accepted local or peer workload transition plus compact GPU telemetry to the scheduler. Queued and running work is counted by destination diff --git a/services/shared/discovery/discovery.go b/services/shared/discovery/discovery.go index 973ac4b2..e5bcd8c1 100644 --- a/services/shared/discovery/discovery.go +++ b/services/shared/discovery/discovery.go @@ -8,10 +8,11 @@ // // The core is a scan-and-diff state machine: each scan browses a service type // over grandcat/zeroconf, re-sends the PTR query from a per-interface unicast -// socket (the Windows send workaround — zeroconf sends from a multicast-bound -// socket Windows refuses to transmit on), and reconciles the result against the -// known-node map. Address/TXT comparison is order-insensitive so a multi-homed -// node whose records come back reordered doesn't churn a spurious "updated". +// UDP 5353 socket (the Windows send workaround — zeroconf sends from a +// multicast-bound socket Windows refuses to transmit on), and reconciles the +// result against the known-node map. Address/TXT comparison is order-insensitive +// so a multi-homed node whose records come back reordered doesn't churn a +// spurious "updated". // // The per-service variations are expressed as functional options rather than // forks: @@ -44,9 +45,10 @@ import ( "sync" "time" + "nvpair-shared/mdns" + "github.com/grandcat/zeroconf" "github.com/miekg/dns" - "golang.org/x/net/ipv4" ) // Event types emitted by Run. @@ -607,8 +609,6 @@ func sendMulticastQuery(service, domain string) map[string]bool { return outcomes } - target := &net.UDPAddr{IP: net.IPv4(224, 0, 0, 251), Port: 5353} - ifaces, err := net.Interfaces() if err != nil { slog.Warn("mdns send: enumerate interfaces failed", "err", err) @@ -630,44 +630,20 @@ func sendMulticastQuery(service, domain string) map[string]bool { if err != nil { continue } - var src net.IP - for _, a := range addrs { - ipnet, ok := a.(*net.IPNet) - if !ok { - continue - } - if ip4 := ipnet.IP.To4(); ip4 != nil { - src = ip4 - break - } - } + ifi := ifi + src, err := sendMulticastQueryOnInterface(buf, &ifi, addrs, mdns.SendFromInterface) if src == nil { continue } - - ifi := ifi - conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: src, Port: 0}) if err != nil { - slog.Debug("mdns send: bind failed", "iface", ifi.Name, "ip", src.String(), "err", err) + slog.Debug("mdns send: send failed", "iface", ifi.Name, "ip", src.String(), "err", err) outcomes[ifi.Name] = false - failures = append(failures, fmt.Sprintf("%s bind: %v", ifi.Name, err)) + failures = append(failures, fmt.Sprintf("%s send: %v", ifi.Name, err)) continue } - pc := ipv4.NewPacketConn(conn) - if err := pc.SetMulticastInterface(&ifi); err != nil { - slog.Debug("mdns send: SetMulticastInterface failed", "iface", ifi.Name, "err", err) - } - _ = pc.SetMulticastTTL(255) - if _, err := conn.WriteToUDP(buf, target); err != nil { - slog.Debug("mdns send: write failed", "iface", ifi.Name, "ip", src.String(), "err", err) - outcomes[ifi.Name] = false - failures = append(failures, fmt.Sprintf("%s write: %v", ifi.Name, err)) - } else { - slog.Debug("mdns send: query sent", "service", service, "iface", ifi.Name, "ip", src.String()) - outcomes[ifi.Name] = true - sent++ - } - _ = conn.Close() + slog.Debug("mdns send: query sent", "service", service, "iface", ifi.Name, "ip", src.String()) + outcomes[ifi.Name] = true + sent++ } if sent == 0 { @@ -681,6 +657,25 @@ func sendMulticastQuery(service, domain string) map[string]bool { return outcomes } +func sendMulticastQueryOnInterface( + buf []byte, + ifi *net.Interface, + addrs []net.Addr, + send func([]byte, *net.Interface, net.IP, *net.UDPAddr) error, +) (net.IP, error) { + for _, addr := range addrs { + ipnet, ok := addr.(*net.IPNet) + if !ok { + continue + } + if ip4 := ipnet.IP.To4(); ip4 != nil { + target := &net.UDPAddr{IP: net.IPv4(224, 0, 0, 251), Port: 5353} + return ip4, send(buf, ifi, ip4, target) + } + } + return nil, nil +} + // UUIDFromTXT returns the value of the "uuid=" TXT record, or "" if absent. It's // the stable per-host identity carried on the node-scanner daemon's single // _nvpair-node record, and was triplicated across the two proxies and the scanner diff --git a/services/shared/discovery/discovery_test.go b/services/shared/discovery/discovery_test.go index 2aaa683e..4b873db0 100644 --- a/services/shared/discovery/discovery_test.go +++ b/services/shared/discovery/discovery_test.go @@ -5,7 +5,9 @@ package discovery import ( "context" + "errors" "fmt" + "net" "sync/atomic" "testing" "time" @@ -446,6 +448,90 @@ func TestRunEmitsAndCloses(t *testing.T) { } } +func TestSendMulticastQueryOnInterfaceUsesFirstIPv4AndMDNSTarget(t *testing.T) { + ifi := &net.Interface{Index: 7, Name: "eth0"} + wantSource := net.IPv4(192, 0, 2, 10) + addrs := []net.Addr{ + &net.IPNet{IP: net.ParseIP("2001:db8::10")}, + &net.IPNet{IP: wantSource}, + &net.IPNet{IP: net.IPv4(198, 51, 100, 20)}, + } + payload := []byte("PTR query") + + var gotPayload []byte + var gotInterface *net.Interface + var gotSource net.IP + var gotTarget *net.UDPAddr + source, err := sendMulticastQueryOnInterface( + payload, + ifi, + addrs, + func(buf []byte, sentIfi *net.Interface, src net.IP, target *net.UDPAddr) error { + gotPayload = append([]byte(nil), buf...) + gotInterface = sentIfi + gotSource = append(net.IP(nil), src...) + gotTarget = target + return nil + }, + ) + if err != nil { + t.Fatalf("sendMulticastQueryOnInterface: %v", err) + } + if !source.Equal(wantSource) || !gotSource.Equal(wantSource) { + t.Fatalf("source = %s / sent %s, want %s", source, gotSource, wantSource) + } + if gotInterface != ifi { + t.Errorf("interface = %v, want %v", gotInterface, ifi) + } + if string(gotPayload) != string(payload) { + t.Errorf("payload = %q, want %q", gotPayload, payload) + } + if gotTarget == nil || !gotTarget.IP.Equal(net.IPv4(224, 0, 0, 251)) || gotTarget.Port != 5353 { + t.Errorf("target = %v, want 224.0.0.251:5353", gotTarget) + } +} + +func TestSendMulticastQueryOnInterfaceReturnsSenderFailure(t *testing.T) { + wantErr := errors.New("send refused") + wantSource := net.IPv4(192, 0, 2, 10) + source, err := sendMulticastQueryOnInterface( + []byte("PTR query"), + &net.Interface{Index: 7, Name: "eth0"}, + []net.Addr{&net.IPNet{IP: wantSource}}, + func([]byte, *net.Interface, net.IP, *net.UDPAddr) error { + return wantErr + }, + ) + if !source.Equal(wantSource) { + t.Fatalf("source = %s, want %s", source, wantSource) + } + if !errors.Is(err, wantErr) { + t.Fatalf("error = %v, want %v", err, wantErr) + } +} + +func TestSendMulticastQueryOnInterfaceSkipsInterfacesWithoutIPv4(t *testing.T) { + called := false + source, err := sendMulticastQueryOnInterface( + []byte("PTR query"), + &net.Interface{Index: 7, Name: "eth0"}, + []net.Addr{&net.IPNet{IP: net.ParseIP("2001:db8::10")}}, + func([]byte, *net.Interface, net.IP, *net.UDPAddr) error { + called = true + return nil + }, + ) + if err != nil { + t.Fatalf("sendMulticastQueryOnInterface: %v", err) + } + if source != nil { + t.Fatalf("source = %s, want nil", source) + } + if called { + t.Fatal("sender called without an IPv4 address") + } +} + // TestSendFailuresNeedARunAndClearOnRecovery: this feeds address selection, so a // single blip must not move a host's canonical address, and one success must undo // the suppression immediately. diff --git a/services/shared/mdns/responder.go b/services/shared/mdns/responder.go index 880b98f1..71c4fca0 100644 --- a/services/shared/mdns/responder.go +++ b/services/shared/mdns/responder.go @@ -16,9 +16,9 @@ // invisible to LAN peers. // // We keep zeroconf's receive trick (join the group on each multicast interface, -// which works fine on Windows) but send every reply/announcement from a -// per-interface unicast-bound socket with SetMulticastInterface set explicitly. -// That path is well-supported on Windows. +// which works fine on Windows) but send every reply/announcement from UDP 5353 +// on a per-interface unicast-bound socket with SetMulticastInterface set +// explicitly. That path is both RFC-compliant and well-supported on Windows. // // This is the single implementation consolidated (the mDNS dedup) from the five // near-identical copies that lived in nvpair-advertiser, @@ -48,7 +48,6 @@ import ( ) const ( - mdnsPort = 5353 // recordTTL matches what zeroconf advertises for non-A records (3200s) // for service-level records, but RFC 6762 §10 says A records SHOULD use // a TTL of 120s to account for IP address changes. We use the shorter @@ -571,9 +570,8 @@ func (r *Responder) sendUnicast(buf []byte, ifIndex int, to net.Addr) { } // sendOnInterface is the core of the Windows send workaround: it transmits buf -// from a fresh unicast-bound socket on the given interface (setting the -// multicast interface + TTL for group targets), never from the multicast-bound -// receive socket that Windows refuses to send from. +// from a fresh UDP 5353 socket bound to the given interface address, never from +// the multicast-bound receive socket that Windows refuses to send from. func (r *Responder) sendOnInterface(buf []byte, ifIndex int, target *net.UDPAddr) error { addrs, ok := r.ifaces()[ifIndex] if !ok || len(addrs) == 0 { @@ -584,19 +582,8 @@ func (r *Responder) sendOnInterface(buf []byte, ifIndex int, target *net.UDPAddr if err != nil { return err } - conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: src, Port: 0}) - if err != nil { - slog.Debug("mdns: bind failed", "iface", ifi.Name, "ip", src.String(), "err", err) - return err - } - defer conn.Close() - if target.IP.IsMulticast() { - pc := ipv4.NewPacketConn(conn) - _ = pc.SetMulticastInterface(ifi) - _ = pc.SetMulticastTTL(255) - } - if _, err := conn.WriteToUDP(buf, target); err != nil { - slog.Debug("mdns: write failed", "iface", ifi.Name, "ip", src.String(), "target", target.String(), "err", err) + if err := SendFromInterface(buf, ifi, src, target); err != nil { + slog.Debug("mdns: send failed", "iface", ifi.Name, "ip", src.String(), "target", target.String(), "err", err) return err } return nil diff --git a/services/shared/mdns/sender.go b/services/shared/mdns/sender.go new file mode 100644 index 00000000..833e4858 --- /dev/null +++ b/services/shared/mdns/sender.go @@ -0,0 +1,102 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package mdns + +import ( + "context" + "errors" + "fmt" + "log/slog" + "net" + + "golang.org/x/net/ipv4" +) + +const ( + mdnsPort = 5353 + preferredMulticastTTL = 255 + fallbackMulticastTTL = 1 +) + +type packetWriter interface { + WriteTo([]byte, net.Addr) (int, error) +} + +type multicastOptions interface { + SetMulticastInterface(*net.Interface) error + SetMulticastTTL(int) error +} + +// SendFromInterface transmits one IPv4 mDNS packet from the selected address +// and the RFC 6762 source port. A fresh source-bound socket preserves reliable +// per-interface egress on Windows while the reuse controls let it coexist with +// the long-lived mDNS receive sockets on every supported platform. Multicast +// interface selection is advisory, and TTL setup falls back explicitly from +// RFC 6762's preferred 255 to link-local 1. Failures are logged but never +// suppress the packet write. +func SendFromInterface(buf []byte, ifi *net.Interface, src net.IP, target *net.UDPAddr) error { + source := src.To4() + if source == nil { + return errors.New("mDNS source is not IPv4") + } + if target == nil || target.IP.To4() == nil || target.Port == 0 { + return errors.New("mDNS target is not a valid IPv4 endpoint") + } + + lc := net.ListenConfig{Control: setSenderReuse} + conn, err := lc.ListenPacket( + context.Background(), + "udp4", + (&net.UDPAddr{IP: source, Port: mdnsPort}).String(), + ) + if err != nil { + return fmt.Errorf("bind mDNS sender: %w", err) + } + defer conn.Close() + + return writePacket(buf, ifi, source, target, conn, ipv4.NewPacketConn(conn)) +} + +func writePacket( + buf []byte, + ifi *net.Interface, + source net.IP, + target *net.UDPAddr, + conn packetWriter, + options multicastOptions, +) error { + if target.IP.IsMulticast() { + if ifi == nil { + return errors.New("mDNS multicast target requires an interface") + } + if err := options.SetMulticastInterface(ifi); err != nil { + slog.Debug("mdns: set multicast interface failed; sending with socket route", + "iface", ifi.Name, + "ip", source.String(), + "target", target.String(), + "err", err) + } + if err := options.SetMulticastTTL(preferredMulticastTTL); err != nil { + slog.Debug("mdns: set multicast TTL failed; retrying with TTL 1", + "iface", ifi.Name, + "ip", source.String(), + "target", target.String(), + "ttl", preferredMulticastTTL, + "err", err) + if fallbackErr := options.SetMulticastTTL(fallbackMulticastTTL); fallbackErr != nil { + slog.Debug("mdns: set multicast fallback TTL failed; sending with socket default", + "iface", ifi.Name, + "ip", source.String(), + "target", target.String(), + "ttl", fallbackMulticastTTL, + "err", fallbackErr) + } + } + } + + if _, err := conn.WriteTo(buf, target); err != nil { + return fmt.Errorf("write mDNS packet: %w", err) + } + return nil +} diff --git a/services/shared/mdns/sender_test.go b/services/shared/mdns/sender_test.go new file mode 100644 index 00000000..fe7b756a --- /dev/null +++ b/services/shared/mdns/sender_test.go @@ -0,0 +1,279 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package mdns + +import ( + "bytes" + "context" + "errors" + "log/slog" + "net" + "slices" + "strings" + "testing" + "time" +) + +type recordingPacketWriter struct { + writes int + payload []byte + target net.Addr + err error +} + +func (w *recordingPacketWriter) WriteTo(payload []byte, target net.Addr) (int, error) { + w.writes++ + w.payload = append([]byte(nil), payload...) + w.target = target + if w.err != nil { + return 0, w.err + } + return len(payload), nil +} + +type recordingMulticastOptions struct { + interfaces []*net.Interface + ttls []int + interfaceErr error + ttlErr error + fallbackTTLErr error +} + +func (o *recordingMulticastOptions) SetMulticastInterface(ifi *net.Interface) error { + o.interfaces = append(o.interfaces, ifi) + return o.interfaceErr +} + +func (o *recordingMulticastOptions) SetMulticastTTL(ttl int) error { + o.ttls = append(o.ttls, ttl) + if ttl == fallbackMulticastTTL { + return o.fallbackTTLErr + } + return o.ttlErr +} + +func captureDebugLogs(t *testing.T) *bytes.Buffer { + t.Helper() + logs := &bytes.Buffer{} + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewTextHandler(logs, &slog.HandlerOptions{Level: slog.LevelDebug}))) + t.Cleanup(func() { slog.SetDefault(previous) }) + return logs +} + +func TestWritePacketConfiguresMulticastAndWritesOnce(t *testing.T) { + ifi := &net.Interface{Index: 7, Name: "eth0"} + source := net.IPv4(192, 0, 2, 10) + target := &net.UDPAddr{IP: net.IPv4(224, 0, 0, 251), Port: mdnsPort} + payload := []byte("multicast payload") + writer := &recordingPacketWriter{} + options := &recordingMulticastOptions{} + + if err := writePacket(payload, ifi, source, target, writer, options); err != nil { + t.Fatalf("writePacket: %v", err) + } + if len(options.interfaces) != 1 || options.interfaces[0] != ifi { + t.Fatalf("multicast interfaces = %v, want [%v]", options.interfaces, ifi) + } + if len(options.ttls) != 1 || options.ttls[0] != preferredMulticastTTL { + t.Fatalf("multicast TTLs = %v, want [255]", options.ttls) + } + if writer.writes != 1 { + t.Fatalf("writes = %d, want 1", writer.writes) + } + if !bytes.Equal(writer.payload, payload) { + t.Errorf("payload = %q, want %q", writer.payload, payload) + } + if writer.target != target { + t.Errorf("target = %v, want %v", writer.target, target) + } +} + +func TestWritePacketSkipsMulticastOptionsForUnicast(t *testing.T) { + source := net.IPv4(127, 0, 0, 1) + target := &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 14318} + writer := &recordingPacketWriter{} + options := &recordingMulticastOptions{} + + if err := writePacket([]byte("unicast payload"), nil, source, target, writer, options); err != nil { + t.Fatalf("writePacket: %v", err) + } + if len(options.interfaces) != 0 || len(options.ttls) != 0 { + t.Fatalf("multicast options used for unicast: interfaces=%v TTLs=%v", options.interfaces, options.ttls) + } + if writer.writes != 1 { + t.Fatalf("writes = %d, want 1", writer.writes) + } +} + +func TestWritePacketReturnsWriteFailure(t *testing.T) { + wantErr := errors.New("send refused") + writer := &recordingPacketWriter{err: wantErr} + source := net.IPv4(127, 0, 0, 1) + target := &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 14318} + + err := writePacket([]byte("unicast payload"), nil, source, target, writer, &recordingMulticastOptions{}) + if !errors.Is(err, wantErr) { + t.Fatalf("error = %v, want wrapped %v", err, wantErr) + } + if writer.writes != 1 { + t.Fatalf("writes = %d, want 1", writer.writes) + } +} + +func TestWritePacketLogsMulticastOptionFailuresAndStillWrites(t *testing.T) { + cases := []struct { + name string + interfaceErr error + ttlErr error + fallbackTTLErr error + wantTTLs []int + wantMessages []string + }{ + { + name: "interface", + interfaceErr: errors.New("interface unavailable"), + wantTTLs: []int{255}, + wantMessages: []string{"set multicast interface failed"}, + }, + { + name: "TTL fallback", + ttlErr: errors.New("TTL unavailable"), + wantTTLs: []int{255, 1}, + wantMessages: []string{"set multicast TTL failed"}, + }, + { + name: "both", + interfaceErr: errors.New("interface unavailable"), + ttlErr: errors.New("TTL unavailable"), + wantTTLs: []int{255, 1}, + wantMessages: []string{"set multicast interface failed", "set multicast TTL failed"}, + }, + { + name: "TTL fallback failure", + ttlErr: errors.New("TTL unavailable"), + fallbackTTLErr: errors.New("fallback TTL unavailable"), + wantTTLs: []int{255, 1}, + wantMessages: []string{"set multicast TTL failed", "set multicast fallback TTL failed"}, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + logs := captureDebugLogs(t) + ifi := &net.Interface{Index: 7, Name: "eth0"} + source := net.IPv4(192, 0, 2, 10) + target := &net.UDPAddr{IP: net.IPv4(224, 0, 0, 251), Port: mdnsPort} + writer := &recordingPacketWriter{} + options := &recordingMulticastOptions{ + interfaceErr: tc.interfaceErr, + ttlErr: tc.ttlErr, + fallbackTTLErr: tc.fallbackTTLErr, + } + + if err := writePacket([]byte("multicast payload"), ifi, source, target, writer, options); err != nil { + t.Fatalf("writePacket: %v", err) + } + if len(options.interfaces) != 1 { + t.Fatalf("interface attempts = %d, want 1", len(options.interfaces)) + } + if !slices.Equal(options.ttls, tc.wantTTLs) { + t.Fatalf("multicast TTLs = %v, want %v", options.ttls, tc.wantTTLs) + } + if writer.writes != 1 { + t.Fatalf("writes = %d, want 1", writer.writes) + } + + gotLogs := logs.String() + for _, message := range tc.wantMessages { + if !strings.Contains(gotLogs, message) { + t.Errorf("logs missing %q:\n%s", message, gotLogs) + } + } + for _, field := range []string{"iface=eth0", "ip=192.0.2.10", "target=224.0.0.251:5353"} { + if !strings.Contains(gotLogs, field) { + t.Errorf("logs missing %q:\n%s", field, gotLogs) + } + } + }) + } +} + +func TestResponderSendUsesMDNSSourcePortAlongsideReceiver(t *testing.T) { + ifi, source := loopbackIPv4(t) + + lc := net.ListenConfig{Control: setReuseAddr} + receiveSocket, err := lc.ListenPacket(context.Background(), "udp4", mdnsTargetV4.String()) + if err != nil { + t.Fatalf("open reusable mDNS receive socket: %v", err) + } + defer receiveSocket.Close() + + sink, err := net.ListenUDP("udp4", &net.UDPAddr{IP: source, Port: 0}) + if err != nil { + t.Fatalf("open UDP sink: %v", err) + } + defer sink.Close() + if err := sink.SetReadDeadline(time.Now().Add(2 * time.Second)); err != nil { + t.Fatalf("set sink deadline: %v", err) + } + + target, ok := sink.LocalAddr().(*net.UDPAddr) + if !ok { + t.Fatalf("sink address has type %T, want *net.UDPAddr", sink.LocalAddr()) + } + responder := &Responder{ + ifaceAddrs: map[int][]net.IP{ + ifi.Index: {source}, + }, + } + payload := []byte("mDNS source-port regression") + if err := responder.sendOnInterface(payload, ifi.Index, target); err != nil { + t.Fatalf("sendOnInterface: %v", err) + } + + buf := make([]byte, len(payload)) + n, from, err := sink.ReadFromUDP(buf) + if err != nil { + t.Fatalf("read UDP sink: %v", err) + } + if !bytes.Equal(buf[:n], payload) { + t.Fatalf("payload = %q, want %q", buf[:n], payload) + } + if !from.IP.Equal(source) { + t.Errorf("source IP = %s, want %s", from.IP, source) + } + if from.Port != mdnsPort { + t.Errorf("source port = %d, want %d", from.Port, mdnsPort) + } +} + +func loopbackIPv4(t *testing.T) (*net.Interface, net.IP) { + t.Helper() + ifaces, err := net.Interfaces() + if err != nil { + t.Fatalf("enumerate interfaces: %v", err) + } + for i := range ifaces { + ifi := &ifaces[i] + if ifi.Flags&net.FlagUp == 0 || ifi.Flags&net.FlagLoopback == 0 { + continue + } + addrs, err := ifi.Addrs() + if err != nil { + continue + } + for _, addr := range addrs { + ipnet, ok := addr.(*net.IPNet) + if !ok { + continue + } + if ip4 := ipnet.IP.To4(); ip4 != nil { + return ifi, ip4 + } + } + } + t.Skip("no up IPv4 loopback interface") + return nil, nil +} diff --git a/services/shared/mdns/senderreuse_darwin.go b/services/shared/mdns/senderreuse_darwin.go new file mode 100644 index 00000000..e2c54ad1 --- /dev/null +++ b/services/shared/mdns/senderreuse_darwin.go @@ -0,0 +1,24 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//go:build darwin + +package mdns + +import "syscall" + +// Darwin requires SO_REUSEPORT in addition to SO_REUSEADDR when a +// unicast-bound sender shares UDP 5353 with the system mDNS responder. Go sets +// the same pair automatically for multicast-address listeners on BSD systems. +func setSenderReuse(network, address string, c syscall.RawConn) error { + if err := setReuseAddr(network, address, c); err != nil { + return err + } + var sockErr error + if err := c.Control(func(fd uintptr) { + sockErr = syscall.SetsockoptInt(int(fd), syscall.SOL_SOCKET, syscall.SO_REUSEPORT, 1) + }); err != nil { + return err + } + return sockErr +} diff --git a/services/shared/mdns/senderreuse_default.go b/services/shared/mdns/senderreuse_default.go new file mode 100644 index 00000000..a531e280 --- /dev/null +++ b/services/shared/mdns/senderreuse_default.go @@ -0,0 +1,12 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//go:build !darwin + +package mdns + +import "syscall" + +func setSenderReuse(network, address string, c syscall.RawConn) error { + return setReuseAddr(network, address, c) +} diff --git a/services/shared/mdns/socketreuse_unix.go b/services/shared/mdns/socketreuse_unix.go index 88d9d56b..8253889d 100644 --- a/services/shared/mdns/socketreuse_unix.go +++ b/services/shared/mdns/socketreuse_unix.go @@ -9,11 +9,11 @@ import "syscall" // setReuseAddr is a net.ListenConfig.Control hook that sets SO_REUSEADDR on // the socket before bind. mDNS requires multiple processes on one host to -// share UDP 5353; SO_REUSEADDR (the same option Go's ListenMulticastUDP and -// grandcat/zeroconf use) lets our responder coexist with the scanner, -// node-info, and any system mDNS responder (Bonjour/Avahi). We intentionally -// do not set SO_REUSEPORT — on Linux it load-balances unicast datagrams -// across the sharing sockets, which would steal unicast mDNS replies. +// share UDP 5353; SO_REUSEADDR lets our responder coexist with the scanner, +// node-info, and any system mDNS responder (Bonjour/Avahi). This base hook +// intentionally omits SO_REUSEPORT because it can load-balance Linux unicast +// datagrams into the wrong socket. Darwin sender sockets add SO_REUSEPORT in +// setSenderReuse, matching that platform's multicast sharing behavior. func setReuseAddr(network, address string, c syscall.RawConn) error { var sockErr error if err := c.Control(func(fd uintptr) {