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