QBONE TUN exchanger async refactor: Pass exchanger an exchanger ;) Final (I think) refactor needed to unblock off-threading TUN writes. Need to pass an exchanger instance to the exchanger for use for underlying-triggered operations. Particularly, if a read packet is link-local and requires a response directly written back to TUN, the passed exchanger will be used instead of `this`. Allows better layering of implementations so when I have an inner "thread-local" exchanger (essentially the current impl) being used by an outer thread-offloading exchanger, the inner exchanger can send writes correctly to the outer exchanger that can write using a worker thread instead of sending writes in-thread. PiperOrigin-RevId: 972116597
diff --git a/quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h b/quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h index 13e96d4..5581d7c 100644 --- a/quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h +++ b/quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h
@@ -8,6 +8,7 @@ #include <cstddef> #include <vector> +#include "absl/base/nullability.h" #include "absl/status/statusor.h" #include "absl/types/span.h" #include "quiche/quic/platform/api/quic_test.h" @@ -29,7 +30,10 @@ (override)); }; - MOCK_METHOD(void, Start, (int read_fd, int write_fd), (override)); + MOCK_METHOD(void, Start, + (int read_fd, int write_fd, + QboneClientPacketExchanger* absl_nullable exchanger), + (override)); MOCK_METHOD(void, Stop, (), (override)); MOCK_METHOD(int, OnReadFromNetworkReady, (int max_packets_to_read), (override));
diff --git a/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h b/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h index 1755d75..7763850 100644 --- a/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h +++ b/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h
@@ -8,9 +8,11 @@ #include <cstddef> #include <vector> +#include "absl/base/nullability.h" #include "absl/status/statusor.h" #include "absl/time/time.h" #include "absl/types/span.h" +#include "quiche/common/quiche_callbacks.h" namespace quic { @@ -40,7 +42,11 @@ // Initializes the exchanger to allow read and write using the given file // descriptors. - virtual void Start(int read_fd, int write_fd) = 0; + // + // If `exchanger` is not null, it will be used instead of `this` to handle any + // underlying-triggered read/write operations. + virtual void Start(int read_fd, int write_fd, + QboneClientPacketExchanger* absl_nullable exchanger) = 0; // Uninitializes the exchanger and blocks until all pending read/write // operations (on- or off-thread) are complete. No Visitor callbacks will be
diff --git a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc index f002ece..010c5d1 100644 --- a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc +++ b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc
@@ -29,6 +29,7 @@ #include "quiche/quic/platform/api/quic_bug_tracker.h" #include "quiche/quic/platform/api/quic_ip_address.h" #include "quiche/quic/platform/api/quic_logging.h" +#include "quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h" #include "quiche/quic/qbone/platform/icmp_packet.h" #include "quiche/quic/qbone/platform/kernel_interface.h" #include "quiche/quic/qbone/platform/netlink_interface.h" @@ -57,18 +58,26 @@ read_fd_ >= 0 || write_fd_ >= 0); } -void TunDevicePacketExchanger::Start(int read_fd, int write_fd) { +void TunDevicePacketExchanger::Start( + int read_fd, int write_fd, + QboneClientPacketExchanger* absl_nullable exchanger) { + if (exchanger == nullptr) { + exchanger = this; + } + // Allow idempotent Start() calls with the same file descriptors, but // otherwise it's a bug to try starting an already started exchanger. QUIC_BUG_IF(qbone_tun_device_packet_exchanger_already_started, - (read_fd_ >= 0 || write_fd_ >= 0) && - (read_fd_ != read_fd || write_fd_ != write_fd)); + (read_fd_ >= 0 || write_fd_ >= 0 || exchanger_ != nullptr) && + (read_fd_ != read_fd || write_fd_ != write_fd || + exchanger_ != exchanger)); QUIC_BUG_IF(qbone_tun_device_packet_exchanger_invalid_read_fd, read_fd < 0); QUIC_BUG_IF(qbone_tun_device_packet_exchanger_invalid_write_fd, write_fd < 0); read_fd_ = read_fd; write_fd_ = write_fd; + exchanger_ = exchanger; } void TunDevicePacketExchanger::Stop() { @@ -76,6 +85,7 @@ // any pending operations to wait for completion. read_fd_ = -1; write_fd_ = -1; + exchanger_ = nullptr; } int TunDevicePacketExchanger::OnReadFromNetworkReady(int max_packets_to_read) { @@ -83,6 +93,9 @@ QUIC_BUG(qbone_tun_device_packet_exchanger_read_with_invalid_fd) << "Invalid file descriptor of the TUN device: " << read_fd_; return 0; + } else if (!exchanger_) { + QUIC_BUG(qbone_tun_device_packet_exchanger_read_with_null_exchanger); + return 0; } int packets_read = 0; @@ -309,9 +322,12 @@ CreateIcmpPacket(ip_hdr->ip6_src, ip_hdr->ip6_src, response_hdr, absl::string_view(payload.get(), payload_size), [this](absl::string_view packet) { - WritePacketToNetwork(absl::MakeSpan( - reinterpret_cast<const std::byte*>(packet.data()), - packet.size())); + QUICHE_DCHECK(exchanger_); + if (exchanger_) { + exchanger_->WritePacketToNetwork(absl::MakeSpan( + reinterpret_cast<const std::byte*>(packet.data()), + packet.size())); + } }); return L2ValidationResult::kValidLinkLocal; }
diff --git a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h index 13d1ed5..2c8130d 100644 --- a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h +++ b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h
@@ -21,6 +21,10 @@ namespace quic { +// Exchanger implementation that does read and write operations synchronously in +// the calling thread, invoking visitor callbacks on that same thread. Safe to +// use separate threads for reading and writing as long as all writes occur on +// the same thread and all reads occur on the same thread. class TunDevicePacketExchanger : public QboneClientPacketExchanger { public: // |mtu| is the mtu of the TUN device. @@ -35,7 +39,8 @@ ~TunDevicePacketExchanger() override; // QboneClientPacketExchanger: - void Start(int read_fd, int write_fd) override; + void Start(int read_fd, int write_fd, + QboneClientPacketExchanger* absl_nullable exchanger) override; void Stop() override; int OnReadFromNetworkReady(int max_packets_to_read) override; void WritePacketToNetwork(absl::Span<const std::byte> packet) override; @@ -61,8 +66,6 @@ L2ValidationResult ValidateL2Headers(const ethhdr& eth_header, absl::Span<const std::byte> packet); - int read_fd_ = -1; - int write_fd_ = -1; KernelInterface* kernel_; NetlinkInterface* netlink_; QboneClientPacketExchanger::Visitor& visitor_; @@ -73,6 +76,11 @@ const bool is_tap_; ethhdr eth_hdr_ = {}; bool eth_hdr_initialized_ = false; + + // -1/nullptr before Start() or after Stop(). + int read_fd_ = -1; + int write_fd_ = -1; + QboneClientPacketExchanger* absl_nullable exchanger_ = nullptr; }; } // namespace quic
diff --git a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger_test.cc b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger_test.cc index 38e7e37..34a3f3d 100644 --- a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger_test.cc +++ b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger_test.cc
@@ -21,7 +21,6 @@ #include "quiche/quic/platform/api/quic_test.h" #include "quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h" #include "quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h" -#include "quiche/quic/qbone/mock_qbone_client.h" #include "quiche/quic/qbone/platform/mock_kernel.h" #include "quiche/quic/qbone/platform/mock_netlink.h" #include "quiche/quic/qbone/platform/netlink_interface.h" @@ -58,10 +57,11 @@ MockKernel mock_kernel_; StrictMock<MockQboneClientPacketExchanger::MockVisitor> mock_visitor_; TunDevicePacketExchanger exchanger_; + StrictMock<MockQboneClientPacketExchanger> delegate_exchanger_; }; TEST_F(TunDevicePacketExchangerTest, WritePacketError) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet = "fake packet"; EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) @@ -83,10 +83,10 @@ } TEST_F(TunDevicePacketExchangerTest, RestartExchanger) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); exchanger_.Stop(); - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet = "fake packet"; EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) @@ -114,7 +114,7 @@ } TEST_F(TunDevicePacketExchangerTest, WritePacketBlocked) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet = "fake packet"; EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) @@ -136,7 +136,7 @@ } TEST_F(TunDevicePacketExchangerTest, WritePacketSuccessfulWrite) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet = "fake packet"; EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) @@ -170,7 +170,7 @@ TunDevicePacketExchanger tap_exchanger(kMtu, &mock_kernel, &mock_netlink, &mock_visitor, /*is_tap=*/true, "tap0"); - tap_exchanger.Start(kReadFd, kWriteFd); + tap_exchanger.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet = "fake packet"; @@ -221,7 +221,7 @@ } TEST_F(TunDevicePacketExchangerTest, ReadPacketError) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) .WillOnce([](int fd, const struct iovec* iov, int iovcnt) { @@ -235,7 +235,7 @@ } TEST_F(TunDevicePacketExchangerTest, ReadPacketBlocked) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) .WillOnce([](int fd, const struct iovec* iov, int iovcnt) { @@ -249,7 +249,7 @@ } TEST_F(TunDevicePacketExchangerTest, ReadPacketSuccessfulRead) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet = "fake_packet"; EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) @@ -271,7 +271,7 @@ } TEST_F(TunDevicePacketExchangerTest, MultipleReadsMoreAvailableThanMax) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet1 = "fake_packet_1"; std::string packet2 = "fake_packet_2"; @@ -305,7 +305,7 @@ } TEST_F(TunDevicePacketExchangerTest, MultipleReadsBlockedBeforeMax) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string packet1 = "fake_packet_1"; @@ -336,7 +336,7 @@ } TEST_F(TunDevicePacketExchangerTest, MultiReadInvalidSizeHuge) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); std::string valid_packet = "valid_packet"; @@ -367,7 +367,7 @@ TEST_F(TunDevicePacketExchangerTest, ReadPacketBlockedOnFirstWithMaxMoreThanOne) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) .WillOnce([](int fd, const struct iovec* iov, int iovcnt) { @@ -392,10 +392,11 @@ StrictMock<MockNetlink> mock_netlink_; StrictMock<MockQboneClientPacketExchanger::MockVisitor> mock_visitor_; TunDevicePacketExchanger exchanger_; + StrictMock<MockQboneClientPacketExchanger> delegate_exchanger_; }; TEST_F(TunDevicePacketExchangerTapTest, ReadPacketTapSuccess) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); ip6_hdr ip_hdr{}; ip_hdr.ip6_vfc = 0x60; // Version 6 @@ -431,7 +432,7 @@ } TEST_F(TunDevicePacketExchangerTapTest, ReadPacketTapInvalidL2) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); ethhdr eth_hdr{}; eth_hdr.h_proto = QuicheEndian::HostToNet16(ETH_P_ARP); // Non-IPv6 @@ -449,7 +450,60 @@ } TEST_F(TunDevicePacketExchangerTapTest, ReadPacketTapNeighborSolicitation) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); + + ip6_hdr ip_hdr{}; + ip_hdr.ip6_vfc = 0x60; // Version 6 + ip_hdr.ip6_nxt = IPPROTO_ICMPV6; + inet_pton(AF_INET6, "fe80::2", &ip_hdr.ip6_src); + inet_pton(AF_INET6, "fe80::1", &ip_hdr.ip6_dst); + + icmp6_hdr icmp_hdr{}; + icmp_hdr.icmp6_type = ND_NEIGHBOR_SOLICIT; + + in6_addr target_address = QboneConstants::GatewayAddress()->GetIPv6(); + + std::string l3_packet = + std::string(reinterpret_cast<char*>(&ip_hdr), sizeof(ip_hdr)) + + std::string(reinterpret_cast<char*>(&icmp_hdr), sizeof(icmp_hdr)) + + std::string(reinterpret_cast<char*>(&target_address), + sizeof(target_address)); + + ethhdr eth_hdr{}; + eth_hdr.h_proto = QuicheEndian::HostToNet16(ETH_P_IPV6); + + EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) + .WillOnce( + [eth_hdr, l3_packet](int fd, const struct iovec* iov, int iovcnt) { + memcpy(iov[0].iov_base, ð_hdr, ETH_HLEN); + memcpy(iov[1].iov_base, l3_packet.data(), l3_packet.size()); + return ETH_HLEN + l3_packet.size(); + }); + + // Expect GetLinkInfo to populate ethhdr on writing neighbor solicit response. + EXPECT_CALL(mock_netlink_, GetLinkInfo("tap0", _)) + .WillOnce( + [](const std::string& ifname, NetlinkInterface::LinkInfo* link_info) { + memset(link_info->hardware_address, 0x12, ETH_ALEN); + return true; + }); + + // Expect neighbor solicitation response to be sent via delegate exchanger. + EXPECT_CALL(delegate_exchanger_, WritePacketToNetwork(_)); + + // OnReadFromNetworkReady should return 1 because packet was handled + // internally (Neighbor Discovery) but still read from network. + EXPECT_EQ(exchanger_.OnReadFromNetworkReady(/*max_packets_to_read=*/1), 1); + + exchanger_.Stop(); +} + +// If not passed a delegate exchanger, TunDevicePacketExchanger should use +// itself to send internal response packets to link-local neighbor solicitation +// requests. +TEST_F(TunDevicePacketExchangerTapTest, + ReadPacketTapNeighborSolicitationSameExchanger) { + exchanger_.Start(kReadFd, kWriteFd, /*exchanger=*/nullptr); ip6_hdr ip_hdr{}; ip_hdr.ip6_vfc = 0x60; // Version 6 @@ -502,7 +556,7 @@ } TEST_F(TunDevicePacketExchangerTapTest, MultiReadInvalidSizeShort) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); ip6_hdr ip_hdr{}; ip_hdr.ip6_vfc = 0x60; // Version 6 @@ -544,7 +598,7 @@ } TEST_F(TunDevicePacketExchangerTapTest, MultiReadInvalidL2) { - exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Start(kReadFd, kWriteFd, &delegate_exchanger_); ethhdr invalid_eth_hdr{}; invalid_eth_hdr.h_proto = QuicheEndian::HostToNet16(ETH_P_ARP); // Non-IPv6