QBONE TUN exchanger async refactor: Make exchanger start/stop more explicit Needed to ensure clean interactions with starting/stopping off-thread writes. Also, now more explicitly a bug to call read/write without a valid socket, rather than being treated as a read/write error, but in a quick look through our code, I don't think we ever try to do read/writes if socket setup fails. Necessary mostly so I can consistently promise the Visitor won't receive callbacks after Stop(). PiperOrigin-RevId: 960557539
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 11db1c1..d42ab2c 100644 --- a/quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h +++ b/quiche/quic/qbone/bonnet/mock_qbone_client_packet_exchanger.h
@@ -29,6 +29,8 @@ (override)); }; + MOCK_METHOD(void, Start, (int read_fd, int write_fd), (override)); + MOCK_METHOD(void, Stop, (), (override)); MOCK_METHOD(bool, ReadAndDeliverPacket, (QboneClientInterface * qbone_client), (override)); MOCK_METHOD(void, WritePacketToNetwork, (const char* packet, size_t size),
diff --git a/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h b/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h index 27862b1..c88b48d 100644 --- a/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h +++ b/quiche/quic/qbone/bonnet/qbone_client_packet_exchanger.h
@@ -39,12 +39,22 @@ virtual ~QboneClientPacketExchanger() = default; + // Initializes the exchanger to allow read and write using the given file + // descriptors. + virtual void Start(int read_fd, int write_fd) = 0; + + // Uninitializes the exchanger and blocks until all pending read/write + // operations (on- or off-thread) are complete. No Visitor callbacks will be + // made after this completes. + virtual void Stop() = 0; + // Reads a packet from the local network and delivers the packet to - // qbone_client. Returns true if there may be more packets to read. + // qbone_client. Returns true if there may be more packets to read. Must not + // be called before Start() or after Stop(). virtual bool ReadAndDeliverPacket(QboneClientInterface* qbone_client) = 0; // Writes a packet to the local network. If the write would be blocked, the - // packet is dropped. + // packet is dropped. Must not be called before Start() or after Stop(). virtual void WritePacketToNetwork(const char* packet, size_t size) = 0; };
diff --git a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc index ccb2770..8489b3e 100644 --- a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc +++ b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.cc
@@ -26,6 +26,7 @@ #include "absl/time/clock.h" #include "absl/time/time.h" #include "absl/types/span.h" +#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/platform/icmp_packet.h" @@ -51,13 +52,37 @@ read_buffer_(mtu), is_tap_(is_tap) {} +TunDevicePacketExchanger::~TunDevicePacketExchanger() { + QUIC_BUG_IF(bonnet_tun_device_packet_exchanger_not_stopped, + read_fd_ >= 0 || write_fd_ >= 0); +} + +void TunDevicePacketExchanger::Start(int read_fd, int write_fd) { + // 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)); + + 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; +} + +void TunDevicePacketExchanger::Stop() { + // This implementation does not employ any worker threads, so there cannot be + // any pending operations to wait for completion. + read_fd_ = -1; + write_fd_ = -1; +} + bool TunDevicePacketExchanger::ReadAndDeliverPacket( QboneClientInterface* qbone_client) { if (read_fd_ < 0) { - absl::Status error = absl::InternalError( - absl::StrCat("Invalid file descriptor of the TUN device: ", read_fd_)); - QUIC_LOG_EVERY_N_SEC(ERROR, 60) << "Packet read failed: " << error; - visitor_.OnRead(std::move(error)); + QUIC_BUG(qbone_tun_device_packet_exchanger_read_with_invalid_fd) + << "Invalid file descriptor of the TUN device: " << read_fd_; return false; } @@ -129,10 +154,8 @@ void TunDevicePacketExchanger::WritePacketToNetwork(const char* packet, size_t size) { if (write_fd_ < 0) { - absl::Status error = absl::InternalError( - absl::StrCat("Invalid file descriptor of the TUN device: ", write_fd_)); - QUIC_LOG_EVERY_N_SEC(ERROR, 60) << "Packet write failed: " << error; - visitor_.OnWrite(std::move(error)); + QUIC_BUG(qbone_tun_device_packet_exchanger_write_with_invalid_fd) + << "Invalid file descriptor of the TUN device: " << write_fd_; return; } @@ -165,13 +188,6 @@ .latency = latency}}); } -void TunDevicePacketExchanger::set_read_file_descriptor(int fd) { - read_fd_ = fd; -} -void TunDevicePacketExchanger::set_write_file_descriptor(int fd) { - write_fd_ = fd; -} - void TunDevicePacketExchanger::InitializeEthHdr() { if (!eth_hdr_initialized_) { NetlinkInterface::LinkInfo link_info{};
diff --git a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h index 0ad7bbf..19b03d5 100644 --- a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h +++ b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger.h
@@ -33,10 +33,11 @@ visitor ABSL_ATTRIBUTE_LIFETIME_BOUND, bool is_tap, absl::string_view ifname); - void set_read_file_descriptor(int fd); - void set_write_file_descriptor(int fd); + ~TunDevicePacketExchanger() override; // QboneClientPacketExchanger: + void Start(int read_fd, int write_fd) override; + void Stop() override; bool ReadAndDeliverPacket(QboneClientInterface* qbone_client) override; void WritePacketToNetwork(const char* packet, size_t size) override;
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 6102d65..cecca0a 100644 --- a/quiche/quic/qbone/bonnet/tun_device_packet_exchanger_test.cc +++ b/quiche/quic/qbone/bonnet/tun_device_packet_exchanger_test.cc
@@ -50,10 +50,7 @@ protected: TunDevicePacketExchangerTest() : exchanger_(kMtu, &mock_kernel_, nullptr, &mock_visitor_, false, - absl::string_view()) { - exchanger_.set_read_file_descriptor(kReadFd); - exchanger_.set_write_file_descriptor(kWriteFd); - } + absl::string_view()) {} ~TunDevicePacketExchangerTest() override = default; @@ -64,6 +61,8 @@ }; TEST_F(TunDevicePacketExchangerTest, WritePacketError) { + exchanger_.Start(kReadFd, kWriteFd); + std::string packet = "fake packet"; EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) .WillOnce([](int fd, const struct iovec* iov, int iovcnt) -> ssize_t { @@ -78,26 +77,16 @@ EXPECT_CALL(mock_visitor_, OnWrite(StatusIs(Ne(absl::StatusCode::kOk)))); exchanger_.WritePacketToNetwork(packet.data(), packet.size()); + + exchanger_.Stop(); } -TEST_F(TunDevicePacketExchangerTest, WritePacketBlocked) { - std::string packet = "fake packet"; - EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) - .WillOnce([](int fd, const struct iovec* iov, int iovcnt) -> ssize_t { - EXPECT_EQ(iov[0].iov_base, nullptr); - EXPECT_EQ(iov[0].iov_len, 0); - EXPECT_THAT(reinterpret_cast<const char*>(iov[1].iov_base), - testing::StrEq("fake packet")); - EXPECT_EQ(iov[1].iov_len, 11); - errno = EAGAIN; - return -1; - }); +TEST_F(TunDevicePacketExchangerTest, RestartExchanger) { + exchanger_.Start(kReadFd, kWriteFd); + exchanger_.Stop(); - EXPECT_CALL(mock_visitor_, OnWrite(StatusIs(Ne(absl::StatusCode::kOk)))); - exchanger_.WritePacketToNetwork(packet.data(), packet.size()); -} + exchanger_.Start(kReadFd, kWriteFd); -TEST_F(TunDevicePacketExchangerTest, WritePacketSuccessfulWrite) { std::string packet = "fake packet"; EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) .WillOnce( @@ -118,6 +107,56 @@ packet.size())))))) .Times(1); exchanger_.WritePacketToNetwork(packet.data(), packet.size()); + + exchanger_.Stop(); +} + +TEST_F(TunDevicePacketExchangerTest, WritePacketBlocked) { + exchanger_.Start(kReadFd, kWriteFd); + + std::string packet = "fake packet"; + EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) + .WillOnce([](int fd, const struct iovec* iov, int iovcnt) -> ssize_t { + EXPECT_EQ(iov[0].iov_base, nullptr); + EXPECT_EQ(iov[0].iov_len, 0); + EXPECT_THAT(reinterpret_cast<const char*>(iov[1].iov_base), + testing::StrEq("fake packet")); + EXPECT_EQ(iov[1].iov_len, 11); + errno = EAGAIN; + return -1; + }); + + EXPECT_CALL(mock_visitor_, OnWrite(StatusIs(Ne(absl::StatusCode::kOk)))); + exchanger_.WritePacketToNetwork(packet.data(), packet.size()); + + exchanger_.Stop(); +} + +TEST_F(TunDevicePacketExchangerTest, WritePacketSuccessfulWrite) { + exchanger_.Start(kReadFd, kWriteFd); + + std::string packet = "fake packet"; + EXPECT_CALL(mock_kernel_, writev(kWriteFd, _, 2)) + .WillOnce( + [&packet](int fd, const struct iovec* iov, int iovcnt) -> ssize_t { + EXPECT_EQ(iov[0].iov_base, nullptr); + EXPECT_EQ(iov[0].iov_len, 0); + EXPECT_THAT(reinterpret_cast<const char*>(iov[1].iov_base), + StrEq(packet)); + EXPECT_EQ(iov[1].iov_len, packet.size()); + return packet.size(); + }); + + EXPECT_CALL( + mock_visitor_, + OnWrite(IsOkAndHolds(ElementsAre(Field( + &QboneClientPacketExchanger::WriteResult::packet, + ElementsAreArray(reinterpret_cast<const std::byte*>(packet.data()), + packet.size())))))) + .Times(1); + exchanger_.WritePacketToNetwork(packet.data(), packet.size()); + + exchanger_.Stop(); } TEST_F(TunDevicePacketExchangerTest, TapWritePacketSuccessful) { @@ -127,7 +166,7 @@ TunDevicePacketExchanger tap_exchanger(kMtu, &mock_kernel, &mock_netlink, &mock_visitor, /*is_tap=*/true, "tap0"); - tap_exchanger.set_write_file_descriptor(kWriteFd); + tap_exchanger.Start(kReadFd, kWriteFd); std::string packet = "fake packet"; @@ -172,9 +211,13 @@ packet.size())))))); tap_exchanger.WritePacketToNetwork(packet.data(), packet.size()); + + tap_exchanger.Stop(); } TEST_F(TunDevicePacketExchangerTest, ReadPacketError) { + exchanger_.Start(kReadFd, kWriteFd); + EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) .WillOnce([](int fd, const struct iovec* iov, int iovcnt) { errno = ECOMM; @@ -182,9 +225,13 @@ }); EXPECT_CALL(mock_visitor_, OnRead(StatusIs(Ne(absl::StatusCode::kOk)))); EXPECT_FALSE(exchanger_.ReadAndDeliverPacket(&mock_client_)); + + exchanger_.Stop(); } TEST_F(TunDevicePacketExchangerTest, ReadPacketBlocked) { + exchanger_.Start(kReadFd, kWriteFd); + EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) .WillOnce([](int fd, const struct iovec* iov, int iovcnt) { errno = EAGAIN; @@ -192,9 +239,13 @@ }); EXPECT_CALL(mock_visitor_, OnRead(StatusIs(Ne(absl::StatusCode::kOk)))); EXPECT_FALSE(exchanger_.ReadAndDeliverPacket(&mock_client_)); + + exchanger_.Stop(); } TEST_F(TunDevicePacketExchangerTest, ReadPacketSuccessfulRead) { + exchanger_.Start(kReadFd, kWriteFd); + std::string packet = "fake_packet"; EXPECT_CALL(mock_kernel_, readv(kReadFd, _, 2)) .WillOnce([packet](int fd, const struct iovec* iov, int iovcnt) { @@ -211,16 +262,15 @@ ElementsAreArray(reinterpret_cast<const std::byte*>(packet.data()), packet.size())))))); EXPECT_TRUE(exchanger_.ReadAndDeliverPacket(&mock_client_)); + + exchanger_.Stop(); } class TunDevicePacketExchangerTapTest : public QuicTest { protected: TunDevicePacketExchangerTapTest() : exchanger_(kMtu, &mock_kernel_, &mock_netlink_, &mock_visitor_, true, - "tap0") { - exchanger_.set_read_file_descriptor(kReadFd); - exchanger_.set_write_file_descriptor(kWriteFd); - } + "tap0") {} ~TunDevicePacketExchangerTapTest() override = default; @@ -232,6 +282,8 @@ }; TEST_F(TunDevicePacketExchangerTapTest, ReadPacketTapSuccess) { + exchanger_.Start(kReadFd, kWriteFd); + ip6_hdr ip_hdr{}; ip_hdr.ip6_vfc = 0x60; // Version 6 ip_hdr.ip6_nxt = 59; // No next header @@ -262,9 +314,13 @@ ElementsAreArray(reinterpret_cast<const std::byte*>(l3_packet.data()), l3_packet.size())))))); EXPECT_TRUE(exchanger_.ReadAndDeliverPacket(&mock_client_)); + + exchanger_.Stop(); } TEST_F(TunDevicePacketExchangerTapTest, ReadPacketTapInvalidL2) { + exchanger_.Start(kReadFd, kWriteFd); + ethhdr eth_hdr{}; eth_hdr.h_proto = QuicheEndian::HostToNet16(ETH_P_ARP); // Non-IPv6 @@ -276,9 +332,13 @@ EXPECT_CALL(mock_visitor_, OnRead(StatusIs(Ne(absl::StatusCode::kOk)))); EXPECT_FALSE(exchanger_.ReadAndDeliverPacket(&mock_client_)); + + exchanger_.Stop(); } TEST_F(TunDevicePacketExchangerTapTest, ReadPacketTapNeighborSolicitation) { + exchanger_.Start(kReadFd, kWriteFd); + ip6_hdr ip_hdr{}; ip_hdr.ip6_vfc = 0x60; // Version 6 ip_hdr.ip6_nxt = IPPROTO_ICMPV6; @@ -325,6 +385,8 @@ // ReadAndDeliverPacket should return false because packet was handled // internally (Neighbor Discovery). EXPECT_FALSE(exchanger_.ReadAndDeliverPacket(&mock_client_)); + + exchanger_.Stop(); } } // namespace