// Copyright 2018 The Chromium Authors
// Use of this source code is governed by a BSD-style license that can be
// found in the LICENSE file.

#include "net/websockets/websocket_basic_stream_adapters.h"

#include <stdint.h>

#include <string>
#include <string_view>
#include <utility>
#include <vector>

#include "base/byte_size.h"
#include "base/check.h"
#include "base/containers/span.h"
#include "base/functional/bind.h"
#include "base/functional/callback.h"
#include "base/memory/raw_ptr.h"
#include "base/memory/scoped_refptr.h"
#include "base/memory/weak_ptr.h"
#include "base/run_loop.h"
#include "base/strings/strcat.h"
#include "base/task/single_thread_task_runner.h"
#include "base/test/bind.h"
#include "base/test/run_until.h"
#include "base/time/default_tick_clock.h"
#include "base/time/time.h"
#include "net/base/host_port_pair.h"
#include "net/base/io_buffer.h"
#include "net/base/ip_address.h"
#include "net/base/ip_endpoint.h"
#include "net/base/net_errors.h"
#include "net/base/network_anonymization_key.h"
#include "net/base/network_handle.h"
#include "net/base/privacy_mode.h"
#include "net/base/proxy_chain.h"
#include "net/base/request_priority.h"
#include "net/base/session_usage.h"
#include "net/base/test_completion_callback.h"
#include "net/cert/cert_verify_result.h"
#include "net/dns/public/host_resolver_results.h"
#include "net/dns/public/secure_dns_policy.h"
#include "net/http/http_network_session.h"
#include "net/http/transport_security_state.h"
#include "net/log/net_log.h"
#include "net/log/net_log_with_source.h"
#include "net/quic/address_utils.h"
#include "net/quic/crypto/proof_verifier_chromium.h"
#include "net/quic/mock_crypto_client_stream_factory.h"
#include "net/quic/mock_quic_data.h"
#include "net/quic/quic_chromium_alarm_factory.h"
#include "net/quic/quic_chromium_client_session.h"
#include "net/quic/quic_chromium_client_session_peer.h"
#include "net/quic/quic_chromium_connection_helper.h"
#include "net/quic/quic_chromium_packet_reader.h"
#include "net/quic/quic_chromium_packet_writer.h"
#include "net/quic/quic_context.h"
#include "net/quic/quic_http_utils.h"
#include "net/quic/quic_server_info.h"
#include "net/quic/quic_session_alias_key.h"
#include "net/quic/quic_session_key.h"
#include "net/quic/quic_test_packet_maker.h"
#include "net/quic/test_quic_crypto_client_config_handle.h"
#include "net/quic/test_task_runner.h"
#include "net/socket/client_socket_handle.h"
#include "net/socket/client_socket_pool.h"
#include "net/socket/next_proto.h"
#include "net/socket/socket_tag.h"
#include "net/socket/socket_test_util.h"
#include "net/socket/stream_socket.h"
#include "net/spdy/spdy_session_key.h"
#include "net/spdy/spdy_test_util_common.h"
#include "net/ssl/ssl_config.h"
#include "net/ssl/ssl_config_service_defaults.h"
#include "net/ssl/ssl_info.h"
#include "net/test/cert_test_util.h"
#include "net/test/gtest_util.h"
#include "net/test/test_data_directory.h"
#include "net/test/test_with_task_environment.h"
#include "net/third_party/quiche/src/quiche/common/http/http_header_block.h"
#include "net/third_party/quiche/src/quiche/common/platform/api/quiche_flags.h"
#include "net/third_party/quiche/src/quiche/common/quiche_buffer_allocator.h"
#include "net/third_party/quiche/src/quiche/common/simple_buffer_allocator.h"
#include "net/third_party/quiche/src/quiche/http2/core/spdy_protocol.h"
#include "net/third_party/quiche/src/quiche/quic/core/crypto/quic_crypto_client_config.h"
#include "net/third_party/quiche/src/quiche/quic/core/frames/quic_blocked_frame.h"
#include "net/third_party/quiche/src/quiche/quic/core/frames/quic_window_update_frame.h"
#include "net/third_party/quiche/src/quiche/quic/core/http/http_constants.h"
#include "net/third_party/quiche/src/quiche/quic/core/http/http_encoder.h"
#include "net/third_party/quiche/src/quiche/quic/core/qpack/qpack_decoder.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_connection.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_connection_id.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_constants.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_error_codes.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_packets.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_time.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_types.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_utils.h"
#include "net/third_party/quiche/src/quiche/quic/core/quic_versions.h"
#include "net/third_party/quiche/src/quiche/quic/platform/api/quic_socket_address.h"
#include "net/third_party/quiche/src/quiche/quic/platform/api/quic_test.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/crypto_test_utils.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/mock_clock.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/mock_connection_id_generator.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/mock_random.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/qpack/qpack_test_utils.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/quic_flow_controller_peer.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/quic_stream_peer.h"
#include "net/third_party/quiche/src/quiche/quic/test_tools/quic_test_utils.h"
#include "net/traffic_annotation/network_traffic_annotation_test_helper.h"
#include "net/websockets/websocket_test_util.h"
#include "testing/gmock/include/gmock/gmock.h"
#include "testing/gtest/include/gtest/gtest.h"
#include "url/gurl.h"
#include "url/scheme_host_port.h"
#include "url/url_constants.h"

namespace net {
class QuicChromiumClientStream;
class SpdySession;
class WebSocketEndpointLockManager;
class X509Certificate;
}  // namespace net

using testing::_;
using testing::AnyNumber;
using testing::Return;
using testing::StrictMock;
using testing::Test;

namespace net::test {

// A helper class that will delete |adapter| when the callback is invoked.
// Used to test that adapters handle being destroyed in their own callbacks.
class DeleterCallback : public TestCompletionCallbackBase {
 public:
  explicit DeleterCallback(
      std::unique_ptr<WebSocketBasicStream::Adapter> adapter)
      : adapter_(std::move(adapter)) {}

  ~DeleterCallback() override = default;

  CompletionOnceCallback callback() {
    return base::BindOnce(&DeleterCallback::OnComplete, base::Unretained(this));
  }

  WebSocketBasicStream::Adapter* adapter() const { return adapter_.get(); }

 private:
  void OnComplete(int result) {
    adapter_.reset();
    SetResult(result);
  }

  std::unique_ptr<WebSocketBasicStream::Adapter> adapter_;
};

class WebSocketClientSocketHandleAdapterTest : public TestWithTaskEnvironment {
 protected:
  WebSocketClientSocketHandleAdapterTest()
      : network_session_(
            SpdySessionDependencies::SpdyCreateSession(&session_deps_)),
        websocket_endpoint_lock_manager_(
            network_session_->websocket_endpoint_lock_manager()) {}

  ~WebSocketClientSocketHandleAdapterTest() override = default;

  bool InitClientSocketHandle(ClientSocketHandle* connection) {
    scoped_refptr<ClientSocketPool::SocketParams> socks_params =
        base::MakeRefCounted<ClientSocketPool::SocketParams>(
            /*allowed_bad_certs=*/std::vector<SSLConfig::CertAndStatus>());
    TestCompletionCallback callback;
    int rv = connection->Init(
        ClientSocketPool::GroupId(
            url::SchemeHostPort(url::kHttpsScheme, "www.example.org", 443),
            PrivacyMode::PRIVACY_MODE_DISABLED, NetworkAnonymizationKey(),
            SecureDnsPolicy::kAllow, /*disable_cert_network_fetches=*/false,
            handles::kInvalidNetworkHandle),
        socks_params, /*proxy_annotation_tag=*/TRAFFIC_ANNOTATION_FOR_TESTS,
        MEDIUM, SocketTag(), ClientSocketPool::RespectLimits::ENABLED,
        callback.callback(), ClientSocketPool::ProxyAuthCallback(),
        network_session_->GetSocketPool(
            HttpNetworkSession::SocketPoolType::kNormal, ProxyChain::Direct()),
        NetLogWithSource());
    rv = callback.GetResult(rv);
    return rv == OK;
  }

  SpdySessionDependencies session_deps_;
  std::unique_ptr<HttpNetworkSession> network_session_;
  raw_ptr<WebSocketEndpointLockManager> websocket_endpoint_lock_manager_;
};

TEST_F(WebSocketClientSocketHandleAdapterTest, Uninitialized) {
  auto connection = std::make_unique<ClientSocketHandle>();
  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_FALSE(adapter.is_initialized());
}

TEST_F(WebSocketClientSocketHandleAdapterTest, IsInitialized) {
  StaticSocketDataProvider data;
  session_deps_.socket_factory->AddSocketDataProvider(&data);
  SSLSocketDataProvider ssl_socket_data(ASYNC, OK);
  session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_socket_data);

  auto connection = std::make_unique<ClientSocketHandle>();
  ClientSocketHandle* const connection_ptr = connection.get();

  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_FALSE(adapter.is_initialized());

  EXPECT_TRUE(InitClientSocketHandle(connection_ptr));

  EXPECT_TRUE(adapter.is_initialized());
}

TEST_F(WebSocketClientSocketHandleAdapterTest, Disconnect) {
  StaticSocketDataProvider data;
  session_deps_.socket_factory->AddSocketDataProvider(&data);
  SSLSocketDataProvider ssl_socket_data(ASYNC, OK);
  session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_socket_data);

  auto connection = std::make_unique<ClientSocketHandle>();
  EXPECT_TRUE(InitClientSocketHandle(connection.get()));

  StreamSocket* const socket = connection->socket();

  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_TRUE(adapter.is_initialized());

  EXPECT_TRUE(socket->IsConnected());
  adapter.Disconnect();
  EXPECT_FALSE(socket->IsConnected());
}

TEST_F(WebSocketClientSocketHandleAdapterTest, Read) {
  MockRead reads[] = {MockRead(SYNCHRONOUS, "foo"), MockRead("bar")};
  StaticSocketDataProvider data(reads, base::span<MockWrite>());
  session_deps_.socket_factory->AddSocketDataProvider(&data);
  SSLSocketDataProvider ssl_socket_data(ASYNC, OK);
  session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_socket_data);

  auto connection = std::make_unique<ClientSocketHandle>();
  EXPECT_TRUE(InitClientSocketHandle(connection.get()));

  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_TRUE(adapter.is_initialized());

  // Buffer larger than each MockRead.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  int rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("foo", std::string_view(read_buf->data(), rv));

  TestCompletionCallback callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(3, rv);
  EXPECT_EQ("bar", std::string_view(read_buf->data(), rv));

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketClientSocketHandleAdapterTest, ReadIntoSmallBuffer) {
  MockRead reads[] = {MockRead(SYNCHRONOUS, "foo"), MockRead("bar")};
  StaticSocketDataProvider data(reads, base::span<MockWrite>());
  session_deps_.socket_factory->AddSocketDataProvider(&data);
  SSLSocketDataProvider ssl_socket_data(ASYNC, OK);
  session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_socket_data);

  auto connection = std::make_unique<ClientSocketHandle>();
  EXPECT_TRUE(InitClientSocketHandle(connection.get()));

  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_TRUE(adapter.is_initialized());

  // Buffer smaller than each MockRead.
  constexpr int kReadBufSize = 2;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  int rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(2, rv);
  EXPECT_EQ("fo", std::string_view(read_buf->data(), rv));

  rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(1, rv);
  EXPECT_EQ("o", std::string_view(read_buf->data(), rv));

  TestCompletionCallback callback1;
  rv = adapter.Read(read_buf.get(), kReadBufSize, callback1.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback1.WaitForResult();
  ASSERT_EQ(2, rv);
  EXPECT_EQ("ba", std::string_view(read_buf->data(), rv));

  rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(1, rv);
  EXPECT_EQ("r", std::string_view(read_buf->data(), rv));

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketClientSocketHandleAdapterTest, Write) {
  MockWrite writes[] = {MockWrite(SYNCHRONOUS, "foo"), MockWrite("bar")};
  StaticSocketDataProvider data(base::span<MockRead>(), writes);
  session_deps_.socket_factory->AddSocketDataProvider(&data);
  SSLSocketDataProvider ssl_socket_data(ASYNC, OK);
  session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_socket_data);

  auto connection = std::make_unique<ClientSocketHandle>();
  EXPECT_TRUE(InitClientSocketHandle(connection.get()));

  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_TRUE(adapter.is_initialized());

  auto write_buf1 = base::MakeRefCounted<StringIOBuffer>("foo");
  int rv =
      adapter.Write(write_buf1.get(), write_buf1->size(),
                    CompletionOnceCallback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_EQ(3, rv);

  auto write_buf2 = base::MakeRefCounted<StringIOBuffer>("bar");
  TestCompletionCallback callback;
  rv = adapter.Write(write_buf2.get(), write_buf2->size(), callback.callback(),
                     TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(3, rv);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

// Test that if both Read() and Write() returns asynchronously,
// the two callbacks are handled correctly.
TEST_F(WebSocketClientSocketHandleAdapterTest, AsyncReadAndWrite) {
  MockRead reads[] = {MockRead("foobar")};
  MockWrite writes[] = {MockWrite("baz")};
  StaticSocketDataProvider data(reads, writes);
  session_deps_.socket_factory->AddSocketDataProvider(&data);
  SSLSocketDataProvider ssl_socket_data(ASYNC, OK);
  session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_socket_data);

  auto connection = std::make_unique<ClientSocketHandle>();
  EXPECT_TRUE(InitClientSocketHandle(connection.get()));

  WebSocketClientSocketHandleAdapter adapter(std::move(connection));
  EXPECT_TRUE(adapter.is_initialized());

  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback read_callback;
  int rv = adapter.Read(read_buf.get(), kReadBufSize, read_callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("baz");
  TestCompletionCallback write_callback;
  rv = adapter.Write(write_buf.get(), write_buf->size(),
                     write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  rv = read_callback.WaitForResult();
  ASSERT_EQ(6, rv);
  EXPECT_EQ("foobar", std::string_view(read_buf->data(), rv));

  rv = write_callback.WaitForResult();
  ASSERT_EQ(3, rv);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

class MockDelegate : public WebSocketSpdyStreamAdapter::Delegate {
 public:
  ~MockDelegate() override = default;
  MOCK_METHOD(void, OnHeadersSent, (), (override));
  MOCK_METHOD(void,
              OnHeadersReceived,
              (const quiche::HttpHeaderBlock&),
              (override));
  MOCK_METHOD(void, OnClose, (int), (override));
};

class WebSocketSpdyStreamAdapterTest : public TestWithTaskEnvironment {
 protected:
  WebSocketSpdyStreamAdapterTest()
      : url_("wss://www.example.org/"),
        key_(HostPortPair::FromURL(url_),
             PRIVACY_MODE_DISABLED,
             ProxyChain::Direct(),
             SessionUsage::kDestination,
             SocketTag(),
             NetworkAnonymizationKey(),
             SecureDnsPolicy::kAllow,
             /*disable_cert_verification_network_fetches=*/false,
             handles::kInvalidNetworkHandle),
        session_(SpdySessionDependencies::SpdyCreateSession(&session_deps_)),
        ssl_(SYNCHRONOUS, OK) {}

  ~WebSocketSpdyStreamAdapterTest() override = default;

  static quiche::HttpHeaderBlock RequestHeaders() {
    return WebSocketHttp2Request("/", "www.example.org:443",
                                 "http://www.example.org", {});
  }

  static quiche::HttpHeaderBlock ResponseHeaders() {
    return WebSocketHttp2Response({});
  }

  void AddSocketData(SocketDataProvider* data) {
    session_deps_.socket_factory->AddSocketDataProvider(data);
  }

  void AddSSLSocketData() {
    ssl_.ssl_info.cert =
        ImportCertFromFile(GetTestCertsDirectory(), "wildcard.pem");
    ASSERT_TRUE(ssl_.ssl_info.cert);
    session_deps_.socket_factory->AddSSLSocketDataProvider(&ssl_);
  }

  base::WeakPtr<SpdySession> CreateSpdySession() {
    return ::net::CreateSpdySession(session_.get(), key_, NetLogWithSource());
  }

  base::WeakPtr<SpdyStream> CreateSpdyStream(
      base::WeakPtr<SpdySession> session) {
    return CreateStreamSynchronously(SPDY_BIDIRECTIONAL_STREAM, session, url_,
                                     LOWEST, NetLogWithSource());
  }

  SpdyTestUtil spdy_util_;
  StrictMock<MockDelegate> mock_delegate_;

 private:
  const GURL url_;
  const SpdySessionKey key_;
  SpdySessionDependencies session_deps_;
  std::unique_ptr<HttpNetworkSession> session_;
  SSLSocketDataProvider ssl_;
};

TEST_F(WebSocketSpdyStreamAdapterTest, Disconnect) {
  MockRead reads[] = {MockRead(ASYNC, ERR_IO_PENDING, 0),
                      MockRead(ASYNC, 0, 1)};
  SequencedSocketData data(reads, base::span<MockWrite>());
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  EXPECT_TRUE(stream);
  adapter.Disconnect();
  EXPECT_FALSE(stream);

  // Read EOF.
  EXPECT_TRUE(session);
  data.Resume();
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest, SendRequestHeadersThenDisconnect) {
  MockRead reads[] = {MockRead(ASYNC, ERR_IO_PENDING, 0),
                      MockRead(ASYNC, 0, 3)};
  spdy::SpdySerializedFrame headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  spdy::SpdySerializedFrame rst(
      spdy_util_.ConstructSpdyRstStream(1, spdy::ERROR_CODE_CANCEL));
  MockWrite writes[] = {CreateMockWrite(headers, 1), CreateMockWrite(rst, 2)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // First read is a pause and it has lower sequence number than first write.
  // Therefore writing headers does not complete while |data| is paused.
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsPaused(); }));

  // Reset the stream before writing completes.
  // OnHeadersSent() will never be called.
  EXPECT_TRUE(stream);
  adapter.Disconnect();
  EXPECT_FALSE(stream);

  // Resume |data|, finish writing headers, and read EOF.
  EXPECT_TRUE(session);
  data.Resume();
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest, OnHeadersSentThenDisconnect) {
  MockRead reads[] = {MockRead(ASYNC, 0, 2)};
  spdy::SpdySerializedFrame headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  spdy::SpdySerializedFrame rst(
      spdy_util_.ConstructSpdyRstStream(1, spdy::ERROR_CODE_CANCEL));
  MockWrite writes[] = {CreateMockWrite(headers, 0), CreateMockWrite(rst, 1)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Finish asynchronous write of headers.  This calls OnHeadersSent().
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  EXPECT_TRUE(stream);
  adapter.Disconnect();
  EXPECT_FALSE(stream);

  // Read EOF.
  EXPECT_TRUE(session);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest, OnHeadersReceivedThenDisconnect) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, 0, 3)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  spdy::SpdySerializedFrame rst(
      spdy_util_.ConstructSpdyRstStream(1, spdy::ERROR_CODE_CANCEL));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(rst, 2)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  EXPECT_TRUE(stream);
  adapter.Disconnect();
  EXPECT_FALSE(stream);

  // Read EOF.
  EXPECT_TRUE(session);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest, ServerClosesConnection) {
  MockRead reads[] = {MockRead(ASYNC, 0, 0)};
  SequencedSocketData data(reads, base::span<MockWrite>());
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest,
       SendRequestHeadersThenServerClosesConnection) {
  MockRead reads[] = {MockRead(ASYNC, 0, 1)};
  spdy::SpdySerializedFrame headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest,
       OnHeadersReceivedThenServerClosesConnection) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, 0, 2)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));
  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

// Previously we failed to detect a half-close by the server that indicated the
// stream should be closed. This test ensures a half-close is correctly
// detected. See https://crbug.com/1151393.
TEST_F(WebSocketSpdyStreamAdapterTest, OnHeadersReceivedThenStreamEnd) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame stream_end(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(stream_end, 2),
                      MockRead(ASYNC, 0, 4)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, /* fin = */ false));
  // The server's END_STREAM must be answered with our own, per RFC 8441
  // section 5, or the server is left in half-closed(local).
  spdy::SpdySerializedFrame client_end_stream(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(client_end_stream, 3)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));
  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  // Must create buffer before `adapter`, since `adapter` doesn't hold onto a
  // reference to it.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  TestCompletionCallback read_callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, read_callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);
  rv = read_callback.WaitForResult();
  EXPECT_EQ(ERR_CONNECTION_CLOSED, rv);
  EXPECT_TRUE(session);
  EXPECT_FALSE(stream);

  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

TEST_F(WebSocketSpdyStreamAdapterTest, DetachDelegate) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, 0, 2)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  // No Delegate methods shall be called after this.
  adapter.DetachDelegate();

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest, Read) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  // First read is the same size as the buffer, next is smaller, last is larger.
  spdy::SpdySerializedFrame data_frame1(
      spdy_util_.ConstructSpdyDataFrame(1, "foo", false));
  spdy::SpdySerializedFrame data_frame2(
      spdy_util_.ConstructSpdyDataFrame(1, "ba", false));
  spdy::SpdySerializedFrame data_frame3(
      spdy_util_.ConstructSpdyDataFrame(1, "rbaz", false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(data_frame1, 2),
                      CreateMockRead(data_frame2, 3),
                      CreateMockRead(data_frame3, 4), MockRead(ASYNC, 0, 5)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  constexpr int kReadBufSize = 3;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(3, rv);
  EXPECT_EQ("foo", std::string_view(read_buf->data(), rv));

  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);

  // HTTP/2 raw_received_bytes() includes the 9-byte frame header for each
  // frame received so far (at least a HEADERS frame and a DATA frame).
  EXPECT_GE(stream->raw_received_bytes().InBytes(),
            spdy::kFrameHeaderSize + rv);

  // Read EOF to destroy the connection and the stream.
  // This calls SpdySession::Delegate::OnClose().
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  // Two socket reads are concatenated by WebSocketSpdyStreamAdapter.
  rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("bar", std::string_view(read_buf->data(), rv));

  rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("baz", std::string_view(read_buf->data(), rv));

  // Even though connection and stream are already closed,
  // WebSocketSpdyStreamAdapter::Delegate::OnClose() is only called after all
  // buffered data are read.
  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

TEST_F(WebSocketSpdyStreamAdapterTest, CallDelegateOnCloseShouldNotCrash) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame data_frame1(
      spdy_util_.ConstructSpdyDataFrame(1, "foo", false));
  spdy::SpdySerializedFrame data_frame2(
      spdy_util_.ConstructSpdyDataFrame(1, "bar", false));
  spdy::SpdySerializedFrame rst(
      spdy_util_.ConstructSpdyRstStream(1, spdy::ERROR_CODE_CANCEL));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(data_frame1, 2),
                      CreateMockRead(data_frame2, 3), CreateMockRead(rst, 4),
                      MockRead(ASYNC, 0, 5)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Buffer larger than each MockRead.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(3, rv);
  EXPECT_EQ("foo", std::string_view(read_buf->data(), rv));

  // Read RST_STREAM to destroy the stream.
  // This calls SpdySession::Delegate::OnClose().
  EXPECT_TRUE(session);
  EXPECT_TRUE(stream);
  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  // Read remaining buffered data.  This will PostTask CallDelegateOnClose().
  rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("bar", std::string_view(read_buf->data(), rv));

  adapter.DetachDelegate();

  // Run CallDelegateOnClose(), which should not crash
  // even if |delegate_| is null.
  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

TEST_F(WebSocketSpdyStreamAdapterTest, Write) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, 0, 3)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  spdy::SpdySerializedFrame data_frame(
      spdy_util_.ConstructSpdyDataFrame(1, "foo", false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(data_frame, 2)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, nullptr, NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("foo");
  TestCompletionCallback callback;
  rv = adapter.Write(write_buf.get(), write_buf->size(), callback.callback(),
                     TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(3, rv);

  // raw_sent_bytes() should include the HEADERS frame + DATA frame sent.
  EXPECT_GE(stream->raw_sent_bytes().InBytes(),
            spdy::kFrameHeaderSize + write_buf->size());

  // Read EOF.
  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

// Test that if both Read() and Write() returns asynchronously,
// the two callbacks are handled correctly.
TEST_F(WebSocketSpdyStreamAdapterTest, AsyncReadAndWrite) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame read_data_frame(
      spdy_util_.ConstructSpdyDataFrame(1, "foobar", false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(read_data_frame, 3),
                      MockRead(ASYNC, 0, 4)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  spdy::SpdySerializedFrame write_data_frame(
      spdy_util_.ConstructSpdyDataFrame(1, "baz", false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(write_data_frame, 2)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, nullptr, NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback read_callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, read_callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("baz");
  TestCompletionCallback write_callback;
  rv = adapter.Write(write_buf.get(), write_buf->size(),
                     write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  rv = read_callback.WaitForResult();
  ASSERT_EQ(6, rv);
  EXPECT_EQ("foobar", std::string_view(read_buf->data(), rv));

  rv = write_callback.WaitForResult();
  ASSERT_EQ(3, rv);

  // Read EOF.
  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

TEST_F(WebSocketSpdyStreamAdapterTest, ReadCallbackDestroysAdapter) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, ERR_IO_PENDING, 2),
                      MockRead(ASYNC, 0, 3)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  auto adapter = std::make_unique<WebSocketSpdyStreamAdapter>(
      stream, &mock_delegate_, NetLogWithSource());
  EXPECT_TRUE(adapter->is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Send headers.
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  WebSocketSpdyStreamAdapter* adapter_raw = adapter.get();
  DeleterCallback callback(std::move(adapter));

  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  rv = adapter_raw->Read(read_buf.get(), kReadBufSize, callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Read EOF while read is pending.  WebSocketSpdyStreamAdapter::OnClose()
  // should not crash if read callback destroys |adapter|.
  data.Resume();
  rv = callback.WaitForResult();
  EXPECT_THAT(rv, IsError(ERR_CONNECTION_CLOSED));

  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest, WriteCallbackDestroysAdapter) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, ERR_IO_PENDING, 2),
                      MockRead(ASYNC, 0, 3)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  auto adapter = std::make_unique<WebSocketSpdyStreamAdapter>(
      stream, &mock_delegate_, NetLogWithSource());
  EXPECT_TRUE(adapter->is_initialized());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Send headers.
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsIdle(); }));

  WebSocketSpdyStreamAdapter* adapter_raw = adapter.get();
  DeleterCallback callback(std::move(adapter));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("foo");
  rv = adapter_raw->Write(write_buf.get(), write_buf->size(),
                          callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Read EOF while write is pending.  WebSocketSpdyStreamAdapter::OnClose()
  // should not crash if write callback destroys |adapter|.
  data.Resume();
  rv = callback.WaitForResult();
  EXPECT_THAT(rv, IsError(ERR_CONNECTION_CLOSED));

  ASSERT_TRUE(base::test::RunUntil([&] { return !session; }));
  EXPECT_FALSE(stream);

  EXPECT_TRUE(data.AllReadDataConsumed());
  EXPECT_TRUE(data.AllWriteDataConsumed());
}

TEST_F(WebSocketSpdyStreamAdapterTest,
       OnCloseOkShouldBeTranslatedToConnectionClose) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame close(
      spdy_util_.ConstructSpdyRstStream(1, spdy::ERROR_CODE_NO_ERROR));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(close, 2), MockRead(ASYNC, 0, 3)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  // Must create buffer before `adapter`, since `adapter` doesn't hold onto a
  // reference to it.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());
  EXPECT_TRUE(adapter.is_initialized());

  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  TestCompletionCallback callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(ERR_CONNECTION_CLOSED, rv);
}

// The send side is closed once the final DATA frame is queued, so a Write()
// racing in after that must fail rather than reach SpdyStream::SendData(). The
// stream outlives the frame, so Write() cannot rely on the stream being gone.
TEST_F(WebSocketSpdyStreamAdapterTest, WriteAfterEndStreamFails) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame stream_end(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(stream_end, 2),
                      MockRead(ASYNC, ERR_IO_PENDING, 3),  // pause here
                      MockRead(ASYNC, 0, 5)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, /* fin = */ false));
  // Sequenced after the pause, so it stays queued and the stream stays alive
  // with its send side already closed.
  spdy::SpdySerializedFrame client_end_stream(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(client_end_stream, 4)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, nullptr, NetLogWithSource());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Receiving END_STREAM queues ours, which cannot be written while paused, so
  // the stream is still alive with its send side closed.
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsPaused(); }));
  ASSERT_TRUE(stream);

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("foo");
  TestCompletionCallback write_callback;
  rv = adapter.Write(write_buf.get(), write_buf->size(),
                     write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_CONNECTION_CLOSED));

  // Rejecting that write must not have cost us the close: once unblocked, the
  // queued END_STREAM still reaches the wire and closes the stream.
  data.Resume();
  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
  EXPECT_FALSE(stream);
}

// If a write is still in flight when the peer's END_STREAM arrives, our own
// END_STREAM is deferred until that write completes. Should the stream be
// destroyed before then, the deferred send must not run against a dead stream.
TEST_F(WebSocketSpdyStreamAdapterTest, StreamClosedWhileEndStreamDeferred) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame stream_end(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, ERR_IO_PENDING, 2),  // pause here
                      CreateMockRead(stream_end, 3),
                      // Kills the session, and with it the stream, while our
                      // END_STREAM is still deferred behind the pending write.
                      MockRead(ASYNC, ERR_CONNECTION_RESET, 4)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, /* fin = */ false));
  // Sequenced after the reads, so it is still in flight when they arrive.
  spdy::SpdySerializedFrame write_data_frame(
      spdy_util_.ConstructSpdyDataFrame(1, "baz", false));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(write_data_frame, 5)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, nullptr, NetLogWithSource());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsPaused(); }));

  // Start a write that cannot complete yet, so END_STREAM must be deferred.
  auto write_buf = base::MakeRefCounted<StringIOBuffer>("baz");
  TestCompletionCallback write_callback;
  rv = adapter.Write(write_buf.get(), write_buf->size(),
                     write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Deliver END_STREAM, which defers our own, and then destroy the stream while
  // that deferred send is still queued.
  data.Resume();
  EXPECT_THAT(write_callback.WaitForResult(), IsError(ERR_CONNECTION_RESET));
  EXPECT_FALSE(stream);
}

// The normal deferred path: a write is still in flight when the peer's
// END_STREAM arrives, so ours waits for it. When that write completes
// successfully it triggers the deferred END_STREAM, which closes the stream and
// reports the close up to a pending read and to the delegate.
TEST_F(WebSocketSpdyStreamAdapterTest,
       DeferredEndStreamSentAfterWriteCompletes) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame stream_end(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, ERR_IO_PENDING, 2),  // pause here
                      CreateMockRead(stream_end, 3), MockRead(ASYNC, 0, 6)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, /* fin = */ false));
  // Sequenced after the END_STREAM read, so it is still in flight when the
  // peer's END_STREAM arrives and ours has to wait for it.
  spdy::SpdySerializedFrame client_data(
      spdy_util_.ConstructSpdyDataFrame(1, "foo", false));
  spdy::SpdySerializedFrame client_end_stream(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(client_data, 4),
                        CreateMockWrite(client_end_stream, 5)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));
  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  // Must create buffer before `adapter`, since `adapter` doesn't hold onto a
  // reference to it.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsPaused(); }));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("foo");
  TestCompletionCallback write_callback;
  rv = adapter.Write(write_buf.get(), write_buf->size(),
                     write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  TestCompletionCallback read_callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, read_callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Deliver END_STREAM, which has to wait for the in-flight write.
  data.Resume();

  // The write still reports success, and only then is our END_STREAM sent.
  EXPECT_EQ(3, write_callback.WaitForResult());

  // Sending it closes the stream, which surfaces as the read completing.
  EXPECT_THAT(read_callback.WaitForResult(), IsError(ERR_CONNECTION_CLOSED));
  EXPECT_FALSE(stream);

  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

// Our END_STREAM is deferred behind a write that is still in flight, and is
// sent before that write's completion callback runs. A write issued from inside
// that callback must therefore be rejected instead of reaching
// SpdyStream::SendData() with the send side already closed.
TEST_F(WebSocketSpdyStreamAdapterTest, ReentrantWriteAfterDeferredEndStream) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  spdy::SpdySerializedFrame stream_end(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      MockRead(ASYNC, ERR_IO_PENDING, 2),  // pause here
                      CreateMockRead(stream_end, 3), MockRead(ASYNC, 0, 6)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, /* fin = */ false));
  // Sequenced after the END_STREAM read, so it is still in flight when the
  // peer's END_STREAM arrives and ours has to be deferred behind it.
  spdy::SpdySerializedFrame client_data(
      spdy_util_.ConstructSpdyDataFrame(1, "foo", false));
  spdy::SpdySerializedFrame client_end_stream(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(client_data, 4),
                        CreateMockWrite(client_end_stream, 5)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, nullptr, NetLogWithSource());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  ASSERT_TRUE(base::test::RunUntil([&] { return data.IsPaused(); }));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("foo");
  auto reentrant_buf = base::MakeRefCounted<StringIOBuffer>("bar");
  TestCompletionCallback unused_callback;
  int write_result = ERR_IO_PENDING;
  int reentrant_result = ERR_IO_PENDING;

  rv = adapter.Write(write_buf.get(), write_buf->size(),
                     base::BindLambdaForTesting([&](int result) {
                       write_result = result;
                       reentrant_result = adapter.Write(
                           reentrant_buf.get(), reentrant_buf->size(),
                           unused_callback.callback(),
                           TRAFFIC_ANNOTATION_FOR_TESTS);
                     }),
                     TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Deliver END_STREAM, which defers ours behind the pending write.
  data.Resume();

  ASSERT_TRUE(
      base::test::RunUntil([&] { return reentrant_result != ERR_IO_PENDING; }));
  EXPECT_EQ(3, write_result);
  EXPECT_THAT(reentrant_result, IsError(ERR_CONNECTION_CLOSED));

  // Both our data frame and the END_STREAM that followed it were written.
  ASSERT_TRUE(
      base::test::RunUntil([&] { return data.AllWriteDataConsumed(); }));
}

// Closing our half of the stream in response to END_STREAM must not discard
// data that arrived but has not been read yet: the stream is gone, but buffered
// data stays readable and Delegate::OnClose() is still deferred until it has
// all been consumed.
TEST_F(WebSocketSpdyStreamAdapterTest, ClosesStreamWithBufferedDataUnread) {
  spdy::SpdySerializedFrame response_headers(
      spdy_util_.ConstructSpdyResponseHeaders(1, ResponseHeaders(), false));
  // Payload and END_STREAM in the same DATA frame, so the stream ends while
  // "bar" is still sitting in the adapter's read queue.
  spdy::SpdySerializedFrame data_frame(
      spdy_util_.ConstructSpdyDataFrame(1, "foobar", true));
  MockRead reads[] = {CreateMockRead(response_headers, 1),
                      CreateMockRead(data_frame, 2), MockRead(ASYNC, 0, 4)};
  spdy::SpdySerializedFrame request_headers(spdy_util_.ConstructSpdyHeaders(
      1, RequestHeaders(), DEFAULT_PRIORITY, /* fin = */ false));
  spdy::SpdySerializedFrame client_end_stream(
      spdy_util_.ConstructSpdyDataFrame(1, "", true));
  MockWrite writes[] = {CreateMockWrite(request_headers, 0),
                        CreateMockWrite(client_end_stream, 3)};
  SequencedSocketData data(reads, writes);
  AddSocketData(&data);
  AddSSLSocketData();

  EXPECT_CALL(mock_delegate_, OnHeadersSent());
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  // Must create buffer before `adapter`, since `adapter` doesn't hold onto a
  // reference to it.
  constexpr int kReadBufSize = 3;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);

  base::WeakPtr<SpdySession> session = CreateSpdySession();
  base::WeakPtr<SpdyStream> stream = CreateSpdyStream(session);
  WebSocketSpdyStreamAdapter adapter(stream, &mock_delegate_,
                                     NetLogWithSource());

  int rv = stream->SendRequestHeaders(RequestHeaders(), MORE_DATA_TO_SEND);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  TestCompletionCallback callback;
  rv = adapter.Read(read_buf.get(), kReadBufSize, callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));
  rv = callback.WaitForResult();
  ASSERT_EQ(3, rv);
  EXPECT_EQ("foo", std::string_view(read_buf->data(), rv));

  // OnClose() is only reported once the buffered data has been drained, so it
  // must not have arrived yet even after the stream is gone.
  EXPECT_CALL(mock_delegate_, OnClose(ERR_CONNECTION_CLOSED));

  // Our END_STREAM is queued when the fin is received and closes the stream
  // once written. The stream therefore outlives the fin by one write.
  ASSERT_TRUE(
      base::test::RunUntil([&] { return data.AllWriteDataConsumed(); }));
  EXPECT_FALSE(stream);

  // The remainder of the payload survived the stream being destroyed.
  rv = adapter.Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("bar", std::string_view(read_buf->data(), rv));

  ASSERT_TRUE(base::test::RunUntil([&] {
    return data.AllReadDataConsumed() && data.AllWriteDataConsumed();
  }));
}

class MockQuicDelegate : public WebSocketQuicStreamAdapter::Delegate {
 public:
  ~MockQuicDelegate() override = default;
  MOCK_METHOD(void, OnHeadersSent, (), (override));
  MOCK_METHOD(void,
              OnHeadersReceived,
              (const quiche::HttpHeaderBlock&),
              (override));
  MOCK_METHOD(void, OnClose, (int), (override));
};

class WebSocketQuicStreamAdapterTest
    : public TestWithTaskEnvironment,
      public ::testing::WithParamInterface<quic::ParsedQuicVersion> {
 protected:
  static quiche::HttpHeaderBlock RequestHeaders() {
    return WebSocketHttp2Request("/", "www.example.org:443",
                                 "http://www.example.org", {});
  }
  WebSocketQuicStreamAdapterTest()
      : version_(GetParam()),
        mock_quic_data_(version_),
        client_data_stream_id1_(quic::QuicUtils::GetFirstBidirectionalStreamId(
            version_.transport_version,
            quic::Perspective::IS_CLIENT)),
        crypto_config_(
            quic::test::crypto_test_utils::ProofVerifierForTesting()),
        connection_id_(quic::test::TestConnectionId(2)),
        client_maker_(version_,
                      connection_id_,
                      &clock_,
                      "mail.example.org",
                      quic::Perspective::IS_CLIENT),
        server_maker_(version_,
                      connection_id_,
                      &clock_,
                      "mail.example.org",
                      quic::Perspective::IS_SERVER),
        peer_addr_(IPAddress(192, 0, 2, 23), 443),
        destination_endpoint_(url::kHttpsScheme, "mail.example.org", 80) {}

  ~WebSocketQuicStreamAdapterTest() override = default;

  void SetUp() override {
    FLAGS_quic_enable_http3_grease_randomness = false;
    clock_.AdvanceTime(quic::QuicTime::Delta::FromMilliseconds(20));
    quic::QuicEnableVersion(version_);
  }

  void TearDown() override {
    EXPECT_TRUE(mock_quic_data_.AllReadDataConsumed());
    EXPECT_TRUE(mock_quic_data_.AllWriteDataConsumed());
  }

  net::QuicChromiumClientSession::Handle* GetQuicSessionHandle() {
    return session_handle_.get();
  }

  // Helper functions for constructing packets sent by the client

  std::unique_ptr<quic::QuicReceivedPacket> ConstructSettingsPacket(
      uint64_t packet_number) {
    return client_maker_.MakeInitialSettingsPacket(packet_number);
  }

  std::unique_ptr<quic::QuicReceivedPacket> ConstructServerDataPacket(
      uint64_t packet_number,
      std::string_view data) {
    quiche::QuicheBuffer buffer = quic::HttpEncoder::SerializeDataFrameHeader(
        data.size(), quiche::SimpleBufferAllocator::Get());
    return server_maker_.Packet(packet_number)
        .AddStreamFrame(
            client_data_stream_id1_, /*fin=*/false,
            base::StrCat(
                {std::string_view(buffer.data(), buffer.size()), data}))
        .Build();
  }

  std::string ConstructDataFrameForVersion(std::string_view body,
                                           quic::ParsedQuicVersion version) {
    DCHECK(version.IsIetfQuic());
    quiche::QuicheBuffer buffer = quic::HttpEncoder::SerializeDataFrameHeader(
        body.size(), quiche::SimpleBufferAllocator::Get());
    return base::StrCat({std::string_view(buffer.data(), buffer.size()), body});
  }

  std::unique_ptr<quic::QuicReceivedPacket> ConstructRstPacket(
      uint64_t packet_number,
      quic::QuicRstStreamErrorCode error_code) {
    return client_maker_.Packet(packet_number)
        .AddStopSendingFrame(client_data_stream_id1_, error_code)
        .AddRstStreamFrame(client_data_stream_id1_, error_code)
        .Build();
  }

  std::unique_ptr<quic::QuicEncryptedPacket> ConstructClientAckPacket(
      uint64_t packet_number,
      uint64_t largest_received,
      uint64_t smallest_received) {
    return client_maker_.Packet(packet_number)
        .AddAckFrame(1, largest_received, smallest_received)
        .Build();
  }

  std::unique_ptr<quic::QuicReceivedPacket> ConstructAckAndRstPacket(
      uint64_t packet_number,
      quic::QuicRstStreamErrorCode error_code,
      uint64_t largest_received,
      uint64_t smallest_received) {
    return client_maker_.Packet(packet_number)
        .AddAckFrame(/*first_received=*/1, largest_received, smallest_received)
        .AddStopSendingFrame(client_data_stream_id1_, error_code)
        .AddRstStreamFrame(client_data_stream_id1_, error_code)
        .Build();
  }

  void Initialize() {
    auto socket = std::make_unique<MockUDPClientSocket>(
        mock_quic_data_.InitializeAndGetSequencedSocketData(), NetLog::Get());
    socket->Connect(peer_addr_);

    runner_ = base::MakeRefCounted<TestTaskRunner>(&clock_);
    helper_ = std::make_unique<QuicChromiumConnectionHelper>(
        &clock_, &random_generator_);
    alarm_factory_ =
        std::make_unique<QuicChromiumAlarmFactory>(runner_.get(), &clock_);
    // Ownership of 'writer' is passed to 'QuicConnection'.
    QuicChromiumPacketWriter* writer = new QuicChromiumPacketWriter(
        socket.get(), base::SingleThreadTaskRunner::GetCurrentDefault().get());
    quic::QuicConnection* connection = new quic::QuicConnection(
        connection_id_, quic::QuicSocketAddress(),
        net::ToQuicSocketAddress(peer_addr_), helper_.get(),
        alarm_factory_.get(), writer, true /* owns_writer */,
        quic::Perspective::IS_CLIENT, quic::test::SupportedVersions(version_),
        connection_id_generator_);
    connection->set_visitor(&visitor_);

    // Load a certificate that is valid for *.example.org
    scoped_refptr<X509Certificate> test_cert(
        ImportCertFromFile(GetTestCertsDirectory(), "wildcard.pem"));
    EXPECT_TRUE(test_cert.get());

    verify_details_.cert_verify_result.verified_cert = test_cert;
    verify_details_.cert_verify_result.is_issued_by_known_root = true;
    crypto_client_stream_factory_.AddProofVerifyDetails(&verify_details_);

    base::TimeTicks dns_end = base::TimeTicks::Now();
    base::TimeTicks dns_start = dns_end - base::Milliseconds(1);

    session_ = std::make_unique<QuicChromiumClientSession>(
        connection, std::move(socket),
        /*stream_factory=*/nullptr, &crypto_client_stream_factory_, &clock_,
        &transport_security_state_, &ssl_config_service_,
        /*server_info=*/nullptr,
        QuicSessionAliasKey(
            url::SchemeHostPort(),
            QuicSessionKey("mail.example.org", 80, PRIVACY_MODE_DISABLED,
                           ProxyChain::Direct(), SessionUsage::kDestination,
                           SocketTag(), NetworkAnonymizationKey(),
                           SecureDnsPolicy::kAllow,
                           /*require_dns_https_alpn=*/false,
                           /*disable_cert_verification_network_fetches=*/false,
                           handles::kInvalidNetworkHandle)),
        /*require_confirmation=*/false,
        /*migrate_session_early_v2=*/false,
        /*migrate_session_on_network_change_v2=*/false,
        /*default_network=*/handles::kInvalidNetworkHandle,
        quic::QuicTime::Delta::FromMilliseconds(
            kDefaultRetransmittableOnWireTimeout.InMilliseconds()),
        /*migrate_idle_session=*/true, /*allow_port_migration=*/false,
        kDefaultIdleSessionMigrationPeriod, /*multi_port_probing_interval=*/0,
        kMaxTimeOnNonDefaultNetwork,
        kMaxMigrationsToNonDefaultNetworkOnWriteError,
        kMaxMigrationsToNonDefaultNetworkOnPathDegrading,
        kQuicYieldAfterPacketsRead,
        quic::QuicTime::Delta::FromMilliseconds(
            kQuicYieldAfterDurationMilliseconds),
        /*cert_verify_flags=*/0, quic::test::DefaultQuicConfig(),
        std::make_unique<TestQuicCryptoClientConfigHandle>(&crypto_config_),
        "CONNECTION_UNKNOWN", dns_start, dns_end,
        /*resolution_details=*/std::nullopt,
        base::DefaultTickClock::GetInstance(),
        base::SingleThreadTaskRunner::GetCurrentDefault().get(),
        /*socket_performance_watcher=*/nullptr, ConnectionEndpointMetadata(),
        /*enable_origin_frame=*/true, /*allow_server_preferred_address=*/true,
        MultiplexedSessionCreationInitiator::kUnknown,
        NetLogWithSource::Make(NetLogSourceType::NONE),
        QuicConnectionReuseDetails());

    session_->Initialize();

    // Enable extended CONNECT protocol (required for WebSocket over HTTP/3).
    session_->OnSetting(quic::SETTINGS_ENABLE_CONNECT_PROTOCOL, 1);

    // Blackhole QPACK decoder stream instead of constructing mock writes.
    session_->qpack_decoder()->set_qpack_stream_sender_delegate(
        &noop_qpack_stream_sender_delegate_);
    TestCompletionCallback callback;
    EXPECT_THAT(session_->CryptoConnect(callback.callback()), IsOk());
    EXPECT_TRUE(session_->OneRttKeysAvailable());
    session_handle_ = session_->CreateHandle(
        url::SchemeHostPort(url::kHttpsScheme, "mail.example.org", 80));
  }

  const quic::ParsedQuicVersion version_;
  MockQuicData mock_quic_data_;
  StrictMock<MockQuicDelegate> mock_delegate_;
  const quic::QuicStreamId client_data_stream_id1_;

 private:
  quic::QuicCryptoClientConfig crypto_config_;
  const quic::QuicConnectionId connection_id_;

 protected:
  QuicTestPacketMaker client_maker_;
  QuicTestPacketMaker server_maker_;
  std::unique_ptr<QuicChromiumClientSession> session_;

 private:
  quic::MockClock clock_;
  std::unique_ptr<QuicChromiumClientSession::Handle> session_handle_;
  scoped_refptr<TestTaskRunner> runner_;
  ProofVerifyDetailsChromium verify_details_;
  MockCryptoClientStreamFactory crypto_client_stream_factory_;
  SSLConfigServiceDefaults ssl_config_service_;
  quic::test::MockConnectionIdGenerator connection_id_generator_;
  std::unique_ptr<QuicChromiumConnectionHelper> helper_;
  std::unique_ptr<QuicChromiumAlarmFactory> alarm_factory_;
  testing::StrictMock<quic::test::MockQuicConnectionVisitor> visitor_;
  TransportSecurityState transport_security_state_;
  IPAddress ip_;
  IPEndPoint peer_addr_;
  quic::test::MockRandom random_generator_{0};
  url::SchemeHostPort destination_endpoint_;
  quic::test::NoopQpackStreamSenderDelegate noop_qpack_stream_sender_delegate_;
};

// Like net::TestCompletionCallback, but for a callback that takes an unbound
// parameter of type WebSocketQuicStreamAdapter.
struct WebSocketQuicStreamAdapterIsPendingHelper {
  bool operator()(
      const std::unique_ptr<WebSocketQuicStreamAdapter>& adapter) const {
    return !adapter;
  }
};

using TestWebSocketQuicStreamAdapterCompletionCallbackBase =
    net::internal::TestCompletionCallbackTemplate<
        std::unique_ptr<WebSocketQuicStreamAdapter>,
        WebSocketQuicStreamAdapterIsPendingHelper>;

class TestWebSocketQuicStreamAdapterCompletionCallback
    : public TestWebSocketQuicStreamAdapterCompletionCallbackBase {
 public:
  base::OnceCallback<void(std::unique_ptr<WebSocketQuicStreamAdapter>)>
  callback();
};

base::OnceCallback<void(std::unique_ptr<WebSocketQuicStreamAdapter>)>
TestWebSocketQuicStreamAdapterCompletionCallback::callback() {
  return base::BindOnce(
      &TestWebSocketQuicStreamAdapterCompletionCallback::SetResult,
      base::Unretained(this));
}

INSTANTIATE_TEST_SUITE_P(QuicVersion,
                         WebSocketQuicStreamAdapterTest,
                         ::testing::ValuesIn(AllSupportedQuicVersions()),
                         ::testing::PrintToStringParamName());

TEST_P(WebSocketQuicStreamAdapterTest, Disconnect) {
  int client_packet_number = 1;

  // Client sends SETTINGS packet during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends RST_STREAM when Disconnect() is called, indicating stream
  // cancellation.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      ConstructRstPacket(client_packet_number++, quic::QUIC_STREAM_CANCELLED));

  // Pause reading initially to control when connection close is processed.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);
  // Simulate server closing the connection.
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  EXPECT_EQ(1u, session_->GetNumActiveStreams());
  adapter->Disconnect();
  EXPECT_EQ(0u, session_->GetNumActiveStreams());

  // `Read()` after `Disconnect()` should return `ERR_UNEXPECTED`.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  EXPECT_EQ(ERR_UNEXPECTED, adapter->Read(read_buf.get(), kReadBufSize,
                                          CompletionOnceCallback()));

  // Session should still be valid after stream disconnect (only the stream is
  // closed, not the connection).
  EXPECT_TRUE(session_->connection()->connected());

  // Read connection close to destroy the session.
  session_->StartReading();
  mock_quic_data_.Resume();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));

  // Session connection should be closed after reading connection close.
  EXPECT_FALSE(session_->connection()->connected());
}

// Tests that a pending read callback is not invoked after `Disconnect()` and is
// safely dropped when the adapter is destroyed.
TEST_P(WebSocketQuicStreamAdapterTest, DisconnectDropsPendingReadCallback) {
  int client_packet_number = 1;
  int server_packet_number = 1;

  // Client sends SETTINGS during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends REQUEST HEADERS.
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          client_packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  // Server sends RESPONSE HEADERS.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 server_packet_number++, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block), nullptr));

  // Pause to let the test issue a `Read()` before any data arrives.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  // Client sends RST_STREAM when `Disconnect()` is called. Must be ASYNC
  // because it follows a read pause.
  mock_quic_data_.AddWrite(
      ASYNC, ConstructAckAndRstPacket(client_packet_number++,
                                      quic::QUIC_STREAM_CANCELLED, 1, 0));

  // Connection close for cleanup.
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  bool headers_received = false;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_)).WillOnce([&]() {
    headers_received = true;
  });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  // Wait for response headers.
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil([&] {
    return headers_received &&
           mock_quic_data_.GetSequencedSocketData()->IsPaused();
  }));

  // `Read()` returns `ERR_IO_PENDING` since no data is available yet.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback read_callback;
  int rv =
      adapter->Read(read_buf.get(), kReadBufSize, read_callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // `Disconnect()` while the read callback is pending.
  adapter->Disconnect();
  EXPECT_FALSE(read_callback.have_result());

  // Destroying the adapter silently drops the pending callback.
  adapter.reset();
  EXPECT_FALSE(read_callback.have_result());

  mock_quic_data_.Resume();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));
  EXPECT_FALSE(read_callback.have_result());
}

// Tests that a pending write callback is not invoked after `Disconnect()` and
// is safely dropped when the adapter is destroyed.
TEST_P(WebSocketQuicStreamAdapterTest, DisconnectDropsPendingWriteCallback) {
  // Set a very low buffer threshold so that buffered data immediately exceeds
  // it, causing Write() to return ERR_IO_PENDING.
  quic::test::QuicFlagSaver flag_saver;
  SetQuicheFlag(quic_buffered_data_threshold, 1);

  int client_packet_number = 1;
  int server_packet_number = 1;

  // Client sends SETTINGS during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends REQUEST HEADERS.
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          client_packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  // Server sends RESPONSE HEADERS.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 server_packet_number++, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block), nullptr));

  // (A) Pause reading so the test can block flow control and issue a `Write()`.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  // When flow control is blocked, QUIC sends a STREAM_BLOCKED frame. The
  // blocked offset is 170 because `WriteHeaders()` already wrote 170 bytes on
  // this stream: a 3-byte HTTP/3 HEADERS frame header and 167 bytes of
  // QPACK-encoded WebSocket CONNECT request headers.
  mock_quic_data_.AddWrite(
      ASYNC, client_maker_.Packet(client_packet_number++)
                 .AddAckFrame(1, 1, 1)
                 .AddFrame(quic::QuicFrame(
                     quic::QuicBlockedFrame(1, client_data_stream_id1_, 170)))
                 .Build());

  // (B) Pause again so the test can call `Disconnect()` before the connection
  // close is delivered.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  // Client sends RST_STREAM when `Disconnect()` is called. Must be ASYNC
  // because it follows a read pause.
  mock_quic_data_.AddWrite(
      ASYNC,
      ConstructRstPacket(client_packet_number++, quic::QUIC_STREAM_CANCELLED));

  // Connection close for cleanup.
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  // Wait for response headers, then hit pause (A).
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return mock_quic_data_.GetSequencedSocketData()->IsPaused(); }));

  // Block flow control to force data to be buffered. Combined with the low
  // buffer threshold (1 byte), this causes CanWriteNewData() to return false
  // and Write() to return ERR_IO_PENDING.
  quic::QuicStream* quic_stream =
      session_->GetActiveStream(client_data_stream_id1_);
  ASSERT_TRUE(quic_stream);
  quic::QuicStreamOffset current_offset =
      quic::test::QuicStreamPeer::SendWindowOffset(quic_stream) -
      quic::test::QuicStreamPeer::SendWindowSize(quic_stream);
  quic::test::QuicStreamPeer::SetSendWindowOffset(quic_stream, current_offset);

  // `Write()` returns `ERR_IO_PENDING` because flow control is blocked and
  // buffered data exceeds the threshold.
  std::string write_data = "foo";
  auto write_buf = base::MakeRefCounted<StringIOBuffer>(write_data);
  TestCompletionCallback write_callback;
  int rv =
      adapter->Write(write_buf.get(), write_buf->size(),
                     write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Resume (A) to let the BLOCKED frame be sent, then hit pause (B). The
  // write callback remains pending because no `WINDOW_UPDATE` is delivered.
  mock_quic_data_.Resume();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return mock_quic_data_.GetSequencedSocketData()->IsPaused(); }));
  EXPECT_FALSE(write_callback.have_result());

  // `Disconnect()` while the write callback is pending.
  adapter->Disconnect();
  EXPECT_FALSE(write_callback.have_result());

  // Destroying the adapter silently drops the pending callback.
  adapter.reset();
  EXPECT_FALSE(write_callback.have_result());

  mock_quic_data_.Resume();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));
  EXPECT_FALSE(write_callback.have_result());
}

TEST_P(WebSocketQuicStreamAdapterTest, AsyncAdapterCreation) {
  constexpr size_t kMaxOpenStreams = 50;

  int packet_number = 1;
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.Packet(packet_number++)
                       .AddStreamsBlockedFrame(/*control_frame_id=*/1,
                                               /*stream_count=*/kMaxOpenStreams,
                                               /*unidirectional=*/false)
                       .Build());

  mock_quic_data_.AddRead(
      ASYNC, server_maker_.Packet(1)
                 .AddMaxStreamsFrame(/*control_frame_id=*/1,
                                     /*stream_count=*/kMaxOpenStreams + 2,
                                     /*unidirectional=*/false)
                 .Build());

  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  Initialize();

  std::vector<QuicChromiumClientStream*> streams;

  for (size_t i = 0; i < kMaxOpenStreams; i++) {
    QuicChromiumClientStream* stream =
        QuicChromiumClientSessionPeer::CreateOutgoingStream(session_.get());
    ASSERT_TRUE(stream);
    streams.push_back(stream);
    EXPECT_EQ(i + 1, session_->GetNumActiveStreams());
  }

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  // Creating an adapter should fail because of the stream limit.
  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_EQ(adapter, nullptr);
  EXPECT_FALSE(callback.have_result());
  EXPECT_EQ(kMaxOpenStreams, session_->GetNumActiveStreams());

  // Read MAX_STREAMS frame that makes it possible to open WebSocket stream.
  session_->StartReading();
  adapter = callback.WaitForResult();
  ASSERT_TRUE(adapter);
  EXPECT_EQ(kMaxOpenStreams + 1, session_->GetNumActiveStreams());

  // Expect OnClose when connection is closed.
  EXPECT_CALL(mock_delegate_, OnClose(_));
  // Close connection.
  mock_quic_data_.Resume();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));
}

TEST_P(WebSocketQuicStreamAdapterTest, SendRequestHeadersThenDisconnect) {
  int packet_number = 1;
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));
  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      ConstructRstPacket(packet_number++, quic::QUIC_STREAM_CANCELLED));

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);
  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  adapter->Disconnect();
}

TEST_P(WebSocketQuicStreamAdapterTest, OnHeadersReceivedThenDisconnect) {
  int packet_number = 1;
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  quiche::HttpHeaderBlock response_header_block =
      WebSocketHttp2Response({{"content-length", "0"}});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 /*packet_number=*/1, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block),
                 /*spdy_headers_frame_length=*/nullptr));
  mock_quic_data_.AddRead(SYNCHRONOUS, ERR_IO_PENDING);
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, ConstructAckAndRstPacket(packet_number++,
                                            quic::QUIC_STREAM_CANCELLED, 1, 0));
  bool headers_received = false;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_))
      .WillOnce([&](const quiche::HttpHeaderBlock& response_headers) {
        auto it = response_headers.find("content-length");
        EXPECT_NE(response_headers.end(), it);
        EXPECT_EQ("0", it->second);
        headers_received = true;
      });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil([&]() { return headers_received; }));

  adapter->Disconnect();

  // After disconnect, the underlying stream is cleared, so byte counts
  // should return 0.
  EXPECT_EQ(0u, adapter->stream_bytes_read());
  EXPECT_EQ(0u, adapter->stream_bytes_written());
}

TEST_P(WebSocketQuicStreamAdapterTest, Read) {
  int packet_number = 1;
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 /*packet_number=*/1, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block),
                 /*spdy_headers_frame_length=*/nullptr));
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  mock_quic_data_.AddRead(ASYNC, ConstructServerDataPacket(2, "foo"));
  mock_quic_data_.AddRead(SYNCHRONOUS,
                          ConstructServerDataPacket(3, "hogehoge"));
  mock_quic_data_.AddRead(SYNCHRONOUS, ERR_IO_PENDING);

  mock_quic_data_.AddWrite(ASYNC,
                           ConstructClientAckPacket(packet_number++, 2, 0));
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, ConstructAckAndRstPacket(packet_number++,
                                            quic::QUIC_STREAM_CANCELLED, 3, 0));

  base::RunLoop run_loop;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_)).WillOnce([&]() {
    run_loop.Quit();
  });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  session_->StartReading();
  run_loop.Run();

  // After receiving response headers, stream_bytes_read() already reflects
  // the HTTP/3 HEADERS frame bytes received at the QUIC stream level.
  uint64_t bytes_after_headers = adapter->stream_bytes_read();
  EXPECT_GT(bytes_after_headers, 0u);

  // Buffer larger than each MockRead.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback read_callback;

  int rv =
      adapter->Read(read_buf.get(), kReadBufSize, read_callback.callback());

  ASSERT_EQ(ERR_IO_PENDING, rv);

  mock_quic_data_.Resume();
  ASSERT_TRUE(
      base::test::RunUntil([&] { return read_callback.have_result(); }));

  rv = read_callback.WaitForResult();
  ASSERT_EQ(3, rv);
  EXPECT_EQ("foo", std::string_view(read_buf->data(), rv));

  // After reading "foo", stream_bytes_read() should have increased by the
  // HTTP/3 DATA frame size: varint frame header + 3-byte payload.
  EXPECT_GT(adapter->stream_bytes_read(), bytes_after_headers);

  rv = adapter->Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(8, rv);
  EXPECT_EQ("hogehoge", std::string_view(read_buf->data(), rv));

  adapter->Disconnect();

  EXPECT_TRUE(mock_quic_data_.AllReadDataConsumed());
  EXPECT_TRUE(mock_quic_data_.AllWriteDataConsumed());
}

TEST_P(WebSocketQuicStreamAdapterTest, ReadIntoSmallBuffer) {
  int packet_number = 1;
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 /*packet_number=*/1, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block),
                 /*spdy_headers_frame_length=*/nullptr));
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);
  // First read is the same size as the buffer, next is smaller, last is larger.
  mock_quic_data_.AddRead(ASYNC, ConstructServerDataPacket(2, "abc"));
  mock_quic_data_.AddRead(SYNCHRONOUS, ConstructServerDataPacket(3, "12"));
  mock_quic_data_.AddRead(SYNCHRONOUS, ConstructServerDataPacket(4, "ABCD"));
  mock_quic_data_.AddRead(SYNCHRONOUS, ERR_IO_PENDING);

  mock_quic_data_.AddWrite(ASYNC,
                           ConstructClientAckPacket(packet_number++, 2, 0));
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, ConstructAckAndRstPacket(packet_number++,
                                            quic::QUIC_STREAM_CANCELLED, 4, 0));

  base::RunLoop run_loop;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_)).WillOnce([&]() {
    run_loop.Quit();
  });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);
  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  session_->StartReading();
  run_loop.Run();

  constexpr int kReadBufSize = 3;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  TestCompletionCallback read_callback;

  int rv =
      adapter->Read(read_buf.get(), kReadBufSize, read_callback.callback());

  ASSERT_EQ(ERR_IO_PENDING, rv);

  mock_quic_data_.Resume();
  ASSERT_TRUE(
      base::test::RunUntil([&] { return read_callback.have_result(); }));

  rv = read_callback.WaitForResult();
  ASSERT_EQ(3, rv);
  EXPECT_EQ("abc", std::string_view(read_buf->data(), rv));

  rv = adapter->Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("12A", std::string_view(read_buf->data(), rv));

  rv = adapter->Read(read_buf.get(), kReadBufSize, CompletionOnceCallback());
  ASSERT_EQ(3, rv);
  EXPECT_EQ("BCD", std::string_view(read_buf->data(), rv));

  adapter->Disconnect();

  EXPECT_TRUE(mock_quic_data_.AllReadDataConsumed());
  EXPECT_TRUE(mock_quic_data_.AllWriteDataConsumed());
}

TEST_P(WebSocketQuicStreamAdapterTest, Write) {
  int client_packet_number = 1;
  int server_packet_number = 1;

  // Client sends SETTINGS
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends REQUEST HEADERS
  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.MakeRequestHeadersPacket(
                       client_packet_number++, client_data_stream_id1_, false,
                       ConvertRequestPriorityToQuicPriority(LOWEST),
                       std::move(request_header_block), nullptr));

  // Server sends RESPONSE HEADERS
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 server_packet_number++, client_data_stream_id1_, false,
                 std::move(response_header_block), nullptr));

  // Server ACKs client packets 1-2 (SETTINGS + REQUEST HEADERS).
  mock_quic_data_.AddRead(ASYNC, server_maker_.Packet(server_packet_number++)
                                     .AddAckFrame(1, 2, 1)
                                     .Build());

  // Client sends ACK + DATA frame in ONE packet.
  // write operation for testing is configured.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.Packet(client_packet_number++)
          .AddAckFrame(1, 2, 1)  // ACK server packets #1 and #2
          .AddStreamFrame(client_data_stream_id1_, false,
                          ConstructDataFrameForVersion("test data", version_))
          .Build());

  // Server ACKs the DATA packet. While the write in this test completed
  // synchronously, processing this ACK is necessary to keep the mock connection
  // state valid and ensure a clean disconnect.
  mock_quic_data_.AddRead(ASYNC,
                          server_maker_.Packet(server_packet_number++)
                              .AddAckFrame(1, 3, 1)
                              .Build());  // Server ACKs client packet 3 (DATA)

  // Client sends RST_STREAM on disconnect.
  // This write happens synchronously when Disconnect() is called.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      ConstructRstPacket(client_packet_number++, quic::QUIC_STREAM_CANCELLED));

  // Add a dummy read to prevent QUIC auto-read crash
  // After ACK at seq 5, QUIC tries to read again - provide EOF/close
  // Must be ASYNC because we're in a callback context (not actively reading)
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  // STARTS TEST EXECUTION

  bool headers_received = false;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_)).WillOnce([&]() {
    headers_received = true;
  });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  // Send request headers (triggers SEQ #0 and #1).
  adapter->WriteHeaders(RequestHeaders(), false);

  // Wait for response headers (triggers SEQ #2).
  session_->StartReading();

  // Wait for server ACKs of headers (SEQ #3 - SETTINGS + REQUEST).
  ASSERT_TRUE(base::test::RunUntil([&] {
    return headers_received &&
           mock_quic_data_.GetSequencedSocketData()->IsIdle();
  }));

  auto write_buf = base::MakeRefCounted<StringIOBuffer>("test data");
  // Before writing data, record stream_bytes_written().
  uint64_t bytes_before_write = adapter->stream_bytes_written();

  // Perform the write. Since the mock socket write is synchronous, this should
  // complete synchronously.
  int rv = adapter->Write(write_buf.get(), write_buf->size(), base::DoNothing(),
                          TRAFFIC_ANNOTATION_FOR_TESTS);

  EXPECT_EQ(9, rv);

  // After writing "test data" (9 bytes), stream_bytes_written() should
  // increase.
  EXPECT_GT(adapter->stream_bytes_written(), bytes_before_write);

  // Run the loop to process pending events, such as the server's ACK for the
  // recently written data. This ensures the mock socket is not in a stopped
  // state when Disconnect() is called.
  ASSERT_TRUE(base::test::RunUntil(
      [&] { return mock_quic_data_.GetSequencedSocketData()->IsIdle(); }));

  // Disconnect (triggers SEQ #6 RST).
  adapter->Disconnect();

  // Process cleanup (SEQ #7).
  ASSERT_TRUE(base::test::RunUntil([&] {
    return mock_quic_data_.AllReadDataConsumed() &&
           mock_quic_data_.AllWriteDataConsumed();
  }));
}

// Tests that the adapter correctly handles being destroyed from within its own
// read callback. This is a regression test to ensure OnClose() doesn't crash
// when the read callback destroys the adapter.
TEST_P(WebSocketQuicStreamAdapterTest, ReadCallbackDestroysAdapter) {
  int packet_number = 1;

  // Client sends SETTINGS during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  // Client sends WebSocket upgrade request headers.
  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  // Server sends response headers accepting the WebSocket upgrade.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 /*packet_number=*/1, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block),
                 /*spdy_headers_frame_length=*/nullptr));

  // Pause here to allow test to set up the read callback before connection
  // close is processed.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  // Server closes the connection unexpectedly (simulates network error or
  // server shutdown).
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  // Create the WebSocket adapter.
  TestWebSocketQuicStreamAdapterCompletionCallback creation_callback;
  auto adapter = session_handle->CreateWebSocketQuicStreamAdapter(
      &mock_delegate_, creation_callback.callback(),
      TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  // Send WebSocket request headers to initiate the handshake.
  adapter->WriteHeaders(RequestHeaders(), false);

  // Start reading from the session and wait until the mock data is paused
  // (after receiving response headers but before connection close).
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return mock_quic_data_.GetSequencedSocketData()->IsPaused(); }));

  // DeleterCallback will destroy the adapter when the read completes.
  DeleterCallback callback(std::move(adapter));

  // Issue an async read. The callback will destroy the adapter when invoked.
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  int rv = callback.adapter()->Read(read_buf.get(), kReadBufSize,
                                    callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Resume the mock data to deliver the connection close event. This triggers
  // the adapter's OnClose(), which invokes the pending read callback. The read
  // callback (DeleterCallback) destroys the adapter, so OnClose() must handle
  // this safely without crashing.
  mock_quic_data_.Resume();
  rv = callback.WaitForResult();

  // The error is ERR_QUIC_PROTOCOL_ERROR because:
  // 1. Socket returns ERR_CONNECTION_CLOSED
  // 2. QuicChromiumPacketReader::OnReadError converts this to
  //    QUIC_PACKET_READ_ERROR
  // 3. WebSocketQuicSpdyStream::MapQuicErrorToNetError maps any non-zero
  //    connection_error() to ERR_QUIC_PROTOCOL_ERROR
  EXPECT_THAT(rv, IsError(ERR_QUIC_PROTOCOL_ERROR));

  // Wait for the QUIC connection to fully close.
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));

  EXPECT_TRUE(mock_quic_data_.AllReadDataConsumed());
  EXPECT_TRUE(mock_quic_data_.AllWriteDataConsumed());
}

// Verifies that WebSocketQuicStreamAdapter::Write() safely returns when
// QuicSpdyStream::WriteOrBufferBody() synchronously closes the QUIC stream and
// a pending read callback deletes the adapter before Write() continues.
TEST_P(WebSocketQuicStreamAdapterTest,
       WriteSyncSocketErrorDestroysAdapterWithPendingRead) {
  int packet_number = 1;

  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.MakeRequestHeadersPacket(
                       packet_number++, client_data_stream_id1_, /*fin=*/false,
                       ConvertRequestPriorityToQuicPriority(LOWEST),
                       RequestHeaders(), nullptr));
  // This write is the DATA packet from WebSocketQuicStreamAdapter::Write().
  // Fail it synchronously to close the connection inside WriteOrBufferBody().
  mock_quic_data_.AddWrite(SYNCHRONOUS, ERR_CONNECTION_REFUSED);

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback creation_callback;
  auto adapter = session_handle->CreateWebSocketQuicStreamAdapter(
      &mock_delegate_, creation_callback.callback(),
      TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  adapter->WriteHeaders(RequestHeaders(), false);

  // Start a Read() that will finish later. DeleterCallback owns the adapter and
  // deletes it when the read callback runs.
  DeleterCallback callback(std::move(adapter));
  constexpr int kReadBufSize = 1024;
  auto read_buf = base::MakeRefCounted<IOBufferWithSize>(kReadBufSize);
  int rv = callback.adapter()->Read(read_buf.get(), kReadBufSize,
                                    callback.callback());
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // This WebSocketQuicStreamAdapter::Write() triggers the mocked socket error.
  // Closing the stream runs the pending read callback, which deletes the
  // adapter. Write() must return without using the deleted adapter.
  auto write_buf = base::MakeRefCounted<StringIOBuffer>("test data");
  rv = callback.adapter()->Write(write_buf.get(), write_buf->size(),
                                 base::DoNothing(),
                                 TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_CONNECTION_CLOSED));
  EXPECT_THAT(callback.WaitForResult(), IsError(ERR_QUIC_PROTOCOL_ERROR));
}

// Tests that the adapter correctly handles being destroyed from within its own
// OnClose() delegate method, when there are no pending read or write
// callbacks.
TEST_P(WebSocketQuicStreamAdapterTest,
       ServerCancelsStreamAndAppDestroysAdapterInOnClose) {
  int packet_number = 1;

  // Client sends SETTINGS during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  // Server cancels the stream by sending STOP_SENDING + RST_STREAM.
  mock_quic_data_.AddRead(ASYNC,
                          server_maker_.Packet(1)
                              .AddStopSendingFrame(client_data_stream_id1_,
                                                   quic::QUIC_STREAM_CANCELLED)
                              .AddRstStreamFrame(client_data_stream_id1_,
                                                 quic::QUIC_STREAM_CANCELLED)
                              .Build());

  // Client responds with RST_STREAM.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           client_maker_.Packet(packet_number++)
                               .AddAckFrame(1, 1, 1)
                               .AddRstStreamFrame(client_data_stream_id1_,
                                                  quic::QUIC_STREAM_CANCELLED)
                               .Build());

  // Pause to ensure OnClose is processed before connection closes.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  std::unique_ptr<WebSocketQuicStreamAdapter> adapter;

  bool on_close_called = false;
  // Expect OnClose and destroy the adapter inside it.
  EXPECT_CALL(mock_delegate_, OnClose(IsError(ERR_ABORTED)))
      .WillOnce([&adapter, &on_close_called](int status) {
        adapter.reset();
        on_close_called = true;
      });

  TestWebSocketQuicStreamAdapterCompletionCallback creation_callback;
  adapter = session_handle->CreateWebSocketQuicStreamAdapter(
      &mock_delegate_, creation_callback.callback(),
      TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  // Start reading from the session. This will process the RST_STREAM from the
  // server and trigger OnClose().
  session_->StartReading();

  ASSERT_TRUE(base::test::RunUntil([&] { return on_close_called; }));
  EXPECT_EQ(nullptr, adapter);

  // Resume the mock data to allow the connection to close.
  mock_quic_data_.Resume();

  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));
}

// Tests that Write returns ERR_IO_PENDING when the stream's send buffer
// exceeds the threshold, and completes asynchronously when buffer space
// becomes available via OnCanWriteNewData.
TEST_P(WebSocketQuicStreamAdapterTest, WritePendingWhenBufferFull) {
  // Set a threshold so we can control when buffer crosses it.
  // With threshold=100, first write of 90 bytes succeeds, second write of 20
  // bytes (total 110 >= 100) returns ERR_IO_PENDING.
  quic::test::QuicFlagSaver flag_saver;
  SetQuicheFlag(quic_buffered_data_threshold, 100);

  int client_packet_number = 1;
  int server_packet_number = 1;

  // Client sends SETTINGS
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends REQUEST HEADERS
  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.MakeRequestHeadersPacket(
                       client_packet_number++, client_data_stream_id1_, false,
                       ConvertRequestPriorityToQuicPriority(LOWEST),
                       std::move(request_header_block), nullptr));

  // Server sends RESPONSE HEADERS
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 server_packet_number++, client_data_stream_id1_, false,
                 std::move(response_header_block), nullptr));

  // Server ACKs client packets 1-2 (SETTINGS + REQUEST HEADERS).
  mock_quic_data_.AddRead(ASYNC, server_maker_.Packet(server_packet_number++)
                                     .AddAckFrame(1, 2, 1)
                                     .Build());

  // First data: 90 bytes (below threshold, will buffer but write succeeds).
  std::string first_data(90, 'x');

  // Second data: 20 bytes (total 110 bytes > 100 threshold, will block).
  std::string second_data(20, 'y');

  // Explicitly block stream/connection flow control by setting SendWindowSize
  // to 0. When attempting to write with flow control = 0, QUIC will buffer the
  // data and send BLOCKED frames to the peer. We expect both stream-level
  // (STREAM_BLOCKED at offset ~170) and connection-level (DATA_BLOCKED at
  // offset ~195) frames, as both were previously restricted. The exact offsets
  // are based on header sizes after the handshake.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.Packet(client_packet_number++)
          .AddAckFrame(1, 2, 1)  // ACK server packets 1-2
          .AddFrame(quic::QuicFrame(quic::QuicBlockedFrame(
              1,                        // control_frame_id
              client_data_stream_id1_,  // stream_id for stream-level BLOCKED
              170)))                    // offset where stream is blocked
          .AddFrame(quic::QuicFrame(quic::QuicBlockedFrame(
              2,  // control_frame_id
              quic::QuicUtils::GetInvalidStreamId(version_.transport_version),
              195)))  // offset where connection is blocked
          .Build());

  // Pause to allow second write before WINDOW_UPDATE arrives
  mock_quic_data_.AddReadPause();

  // Server sends WINDOW_UPDATE frames (MAX_STREAM_DATA and MAX_DATA in QUIC)
  // to unblock both stream and connection flow control. The large offset
  // (10000) ensures enough credit for all buffered data (headers ~170 bytes +
  // ~115 bytes buffered data).
  mock_quic_data_.AddRead(
      ASYNC,
      server_maker_.Packet(server_packet_number++)
          .AddFrame(quic::QuicFrame(quic::QuicWindowUpdateFrame(
              1,                        // control_frame_id
              client_data_stream_id1_,  // stream_id for stream-level
              10000)))                  // max_data: large enough for all data
          .AddFrame(quic::QuicFrame(quic::QuicWindowUpdateFrame(
              2,  // control_frame_id
              quic::QuicUtils::GetInvalidStreamId(version_.transport_version),
              10000)))  // max_data for connection
          .Build());

  // When flow control unblocks, QUIC will flush ALL buffered data in a single
  // packet. Both first_data (90 bytes) and second_data (20 bytes) were queued
  // while flow control was blocked, so they get sent together as one coalesced
  // stream frame.
  std::string combined_data =
      ConstructDataFrameForVersion(first_data, version_) +
      ConstructDataFrameForVersion(second_data, version_);

  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.Packet(client_packet_number++)
          .AddAckFrame(1, 3, 1)  // ACK server packets 1-3
          .AddStreamFrame(client_data_stream_id1_, false, combined_data)
          .Build());

  // Server ACKs the client packets (SETTINGS, HEADERS, BLOCKED, combined DATA).
  mock_quic_data_.AddRead(ASYNC, server_maker_.Packet(server_packet_number++)
                                     .AddAckFrame(1, 4, 1)
                                     .Build());

  // Client sends RST_STREAM on disconnect.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      ConstructRstPacket(client_packet_number++, quic::QUIC_STREAM_CANCELLED));

  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  // STARTS TEST EXECUTION

  bool headers_received = false;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_)).WillOnce([&]() {
    headers_received = true;
  });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  // Wait for response headers and server ACK.
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil([&] {
    return headers_received &&
           mock_quic_data_.GetSequencedSocketData()->IsIdle();
  }));

  // Restrict the stream's flow control send window to prevent any data
  // from being sent immediately. We set the offset to current bytes_sent,
  // which makes SendWindowSize = 0 (blocked).
  // This forces all data to accumulate in the buffer, allowing us to test
  // the threshold-based blocking behavior.
  quic::QuicStream* quic_stream =
      session_->GetActiveStream(client_data_stream_id1_);
  ASSERT_TRUE(quic_stream);

  // Get the current bytes sent so we can set window to exactly that offset
  // (making remaining window = 0)
  quic::QuicStreamOffset current_stream_offset =
      quic::test::QuicStreamPeer::SendWindowOffset(quic_stream) -
      quic::test::QuicStreamPeer::SendWindowSize(quic_stream);

  // Set stream flow control: offset = current bytes sent, so window size = 0
  quic::test::QuicStreamPeer::SetSendWindowOffset(quic_stream,
                                                  current_stream_offset);

  // Also restrict the connection-level flow control similarly
  quic::QuicStreamOffset current_conn_offset =
      quic::test::QuicFlowControllerPeer::SendWindowOffset(
          session_->flow_controller()) -
      quic::test::QuicFlowControllerPeer::SendWindowSize(
          session_->flow_controller());
  quic::test::QuicFlowControllerPeer::SetSendWindowOffset(
      session_->flow_controller(), current_conn_offset);

  // First write: 90 bytes. With flow control blocked, all data is buffered.
  // BufferedDataBytes = 90 < threshold(100), so CanWriteNewData() = true.
  // Write should complete synchronously with return value = 90.
  auto write_buf1 = base::MakeRefCounted<StringIOBuffer>(first_data);
  int rv = adapter->Write(write_buf1.get(), write_buf1->size(),
                          base::DoNothing(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_EQ(static_cast<int>(first_data.size()), rv);

  // Second write: 20 bytes. Total buffered = 90 + 20 = 110 bytes.
  // BufferedDataBytes = 110 >= threshold(100), so CanWriteNewData() = false.
  // Write should return ERR_IO_PENDING.
  auto write_buf2 = base::MakeRefCounted<StringIOBuffer>(second_data);
  TestCompletionCallback write_callback2;
  rv = adapter->Write(write_buf2.get(), write_buf2->size(),
                      write_callback2.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_EQ(ERR_IO_PENDING, rv);

  // Resume the mock data - this will deliver the MAX_STREAM_DATA frame,
  // which opens flow control window, triggers sending of buffered data,
  // calls OnCanWriteNewData, and completes the second write.
  mock_quic_data_.Resume();

  // Wait for the second write to complete.
  rv = write_callback2.WaitForResult();
  EXPECT_EQ(static_cast<int>(second_data.size()), rv);

  // Wait for the server to ACK the combined DATA packet (server packet #4).
  ASSERT_TRUE(base::test::RunUntil(
      [&] { return mock_quic_data_.GetSequencedSocketData()->IsIdle(); }));

  adapter->Disconnect();

  // Process cleanup.
  ASSERT_TRUE(base::test::RunUntil([&] {
    return mock_quic_data_.AllReadDataConsumed() &&
           mock_quic_data_.AllWriteDataConsumed();
  }));
}

// Tests that receiving a RST_STREAM from the server while a Write() is pending
// correctly completes the write callback with an error.
TEST_P(WebSocketQuicStreamAdapterTest, RstStreamReceivedWhileWritePending) {
  quic::test::QuicFlagSaver flag_saver;
  SetQuicheFlag(quic_buffered_data_threshold, 100);

  int client_packet_number = 1;
  int server_packet_number = 1;

  // Client sends SETTINGS.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends REQUEST HEADERS.
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.MakeRequestHeadersPacket(
                       client_packet_number++, client_data_stream_id1_, false,
                       ConvertRequestPriorityToQuicPriority(LOWEST),
                       std::move(request_header_block), nullptr));

  // Server sends RESPONSE HEADERS.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 server_packet_number++, client_data_stream_id1_, false,
                 std::move(response_header_block), nullptr));

  // Server ACKs client packets 1-2 (SETTINGS + REQUEST HEADERS).
  mock_quic_data_.AddRead(ASYNC, server_maker_.Packet(server_packet_number++)
                                     .AddAckFrame(1, 2, 1)
                                     .Build());

  // Client sends ACK + BLOCKED frames when flow control is exhausted.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.Packet(client_packet_number++)
          .AddAckFrame(1, 2, 1)
          .AddFrame(quic::QuicFrame(
              quic::QuicBlockedFrame(1, client_data_stream_id1_, 170)))
          .AddFrame(quic::QuicFrame(quic::QuicBlockedFrame(
              2,
              quic::QuicUtils::GetInvalidStreamId(version_.transport_version),
              195)))
          .Build());

  mock_quic_data_.AddReadPause();

  // Server terminates stream with STOP_SENDING + RST_STREAM.
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.Packet(server_packet_number++)
                 .AddStopSendingFrame(client_data_stream_id1_,
                                      quic::QUIC_STREAM_PEER_GOING_AWAY)
                 .AddRstStreamFrame(client_data_stream_id1_,
                                    quic::QUIC_STREAM_PEER_GOING_AWAY)
                 .Build());

  // Client MUST respond with RST_STREAM.
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.Packet(client_packet_number++)
                       .AddAckFrame(1, 3, 1)
                       .AddRstStreamFrame(client_data_stream_id1_,
                                          quic::QUIC_STREAM_PEER_GOING_AWAY)
                       .Build());

  // Session continues reading after stream closes. Required for test teardown.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  bool headers_received = false;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_)).WillOnce([&]() {
    headers_received = true;
  });
  EXPECT_CALL(mock_delegate_, OnClose(_));

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  // Wait for response headers and server ACK.
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil([&] {
    return headers_received &&
           mock_quic_data_.GetSequencedSocketData()->IsIdle();
  }));

  // Block flow control to force writes to buffer.
  quic::QuicStream* quic_stream =
      session_->GetActiveStream(client_data_stream_id1_);
  ASSERT_TRUE(quic_stream);

  quic::QuicStreamOffset current_stream_offset =
      quic::test::QuicStreamPeer::SendWindowOffset(quic_stream) -
      quic::test::QuicStreamPeer::SendWindowSize(quic_stream);
  quic::test::QuicStreamPeer::SetSendWindowOffset(quic_stream,
                                                  current_stream_offset);

  quic::QuicStreamOffset current_conn_offset =
      quic::test::QuicFlowControllerPeer::SendWindowOffset(
          session_->flow_controller()) -
      quic::test::QuicFlowControllerPeer::SendWindowSize(
          session_->flow_controller());
  quic::test::QuicFlowControllerPeer::SetSendWindowOffset(
      session_->flow_controller(), current_conn_offset);

  // First write (90 bytes): buffers, returns synchronously.
  std::string first_data(90, 'x');
  auto write_buf1 = base::MakeRefCounted<StringIOBuffer>(first_data);
  int rv = adapter->Write(write_buf1.get(), write_buf1->size(),
                          base::DoNothing(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_EQ(static_cast<int>(first_data.size()), rv);

  // Second write (20 bytes): total 110 >= threshold, returns ERR_IO_PENDING.
  std::string second_data(20, 'y');
  auto write_buf2 = base::MakeRefCounted<StringIOBuffer>(second_data);
  TestCompletionCallback write_callback;
  rv = adapter->Write(write_buf2.get(), write_buf2->size(),
                      write_callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_EQ(ERR_IO_PENDING, rv);

  // Resume mock data to deliver the RST_STREAM while write is pending.
  mock_quic_data_.Resume();

  // Write callback completes with error.
  rv = write_callback.WaitForResult();
  EXPECT_THAT(rv, IsError(ERR_QUIC_PROTOCOL_ERROR));

  // Connection session cleanup.
  mock_quic_data_.Resume();
  ASSERT_TRUE(base::test::RunUntil([&] {
    return mock_quic_data_.AllReadDataConsumed() &&
           mock_quic_data_.AllWriteDataConsumed();
  }));
}

// Tests that OnClose is called when the server sends a CONNECTION_CLOSE frame.
TEST_P(WebSocketQuicStreamAdapterTest, ConnectionCloseTriggersOnClose) {
  int client_packet_number = 1;
  int server_packet_number = 1;

  // Client sends SETTINGS.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(client_packet_number++));

  // Client sends REQUEST HEADERS.
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS, client_maker_.MakeRequestHeadersPacket(
                       client_packet_number++, client_data_stream_id1_, false,
                       ConvertRequestPriorityToQuicPriority(LOWEST),
                       std::move(request_header_block), nullptr));

  // Server sends RESPONSE HEADERS.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 server_packet_number++, client_data_stream_id1_, false,
                 std::move(response_header_block), nullptr));

  // Server sends CONNECTION_CLOSE frame (connection-level error).
  // After receiving CONNECTION_CLOSE, the session is closed and no more
  // reads occur, so we don't add additional mock reads.
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.Packet(server_packet_number++)
                 .AddConnectionCloseFrame(quic::QUIC_PEER_GOING_AWAY,
                                          "Server shutting down")
                 .Build());

  bool on_close_called = false;
  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));
  EXPECT_CALL(mock_delegate_, OnClose(ERR_QUIC_PROTOCOL_ERROR)).WillOnce([&]() {
    on_close_called = true;
  });

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  TestWebSocketQuicStreamAdapterCompletionCallback callback;
  std::unique_ptr<WebSocketQuicStreamAdapter> adapter =
      session_handle->CreateWebSocketQuicStreamAdapter(
          &mock_delegate_, callback.callback(), TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  adapter->WriteHeaders(RequestHeaders(), false);

  // Start reading to process server packets. This will receive the response
  // headers and then the CONNECTION_CLOSE frame, triggering OnClose.
  session_->StartReading();

  ASSERT_TRUE(base::test::RunUntil([&] { return on_close_called; }));
}

// Tests that the adapter correctly handles being destroyed from within its own
// write callback when a socket write error occurs. This covers the scenario
// where network transmission fails (e.g., network disconnected, ICMP
// unreachable).
TEST_P(WebSocketQuicStreamAdapterTest,
       WriteCallbackDestroysAdapterOnSocketWriteError) {
  // Set a very low buffer threshold. When combined with flow control blocking,
  // any buffered data will exceed this threshold and cause Write() to return
  // ERR_IO_PENDING.
  quic::test::QuicFlagSaver flag_saver;
  SetQuicheFlag(quic_buffered_data_threshold, 1);

  int packet_number = 1;

  // Client sends SETTINGS during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  // Client sends WebSocket upgrade request headers.
  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  // Server sends response headers accepting the WebSocket upgrade.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 /*packet_number=*/1, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block),
                 /*spdy_headers_frame_length=*/nullptr));

  // Pause here to allow test to block flow control and issue a Write() before
  // processing further mock data.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  // When flow control is blocked and we try to write, QUIC attempts to send
  // a BLOCKED frame. This socket write FAILS with ERR_FAILED, simulating
  // a network transmission error. This triggers QUIC_PACKET_WRITE_ERROR.
  mock_quic_data_.AddWrite(ASYNC, ERR_FAILED);

  // After write error, connection is closed. Provide EOF for cleanup.
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  // Create the WebSocket adapter.
  TestWebSocketQuicStreamAdapterCompletionCallback creation_callback;
  auto adapter = session_handle->CreateWebSocketQuicStreamAdapter(
      &mock_delegate_, creation_callback.callback(),
      TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  // Send WebSocket request headers to initiate the handshake.
  adapter->WriteHeaders(RequestHeaders(), false);

  // Start reading from the session and wait until the mock data is paused
  // (after receiving response headers but before connection close).
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return mock_quic_data_.GetSequencedSocketData()->IsPaused(); }));

  // Block flow control to force data to be buffered. Combined with the low
  // buffer threshold (1 byte), this causes CanWriteNewData() to return false
  // and Write() to return ERR_IO_PENDING.
  quic::QuicStream* quic_stream =
      session_->GetActiveStream(client_data_stream_id1_);
  ASSERT_TRUE(quic_stream);
  quic::QuicStreamOffset current_offset =
      quic::test::QuicStreamPeer::SendWindowOffset(quic_stream) -
      quic::test::QuicStreamPeer::SendWindowSize(quic_stream);
  quic::test::QuicStreamPeer::SetSendWindowOffset(quic_stream, current_offset);

  // DeleterCallback will destroy the adapter when the write completes.
  DeleterCallback callback(std::move(adapter));

  // Issue an async write. The callback will destroy the adapter when invoked.
  // Write returns ERR_IO_PENDING because:
  // 1. Flow control is blocked, so data is buffered (not sent)
  // 2. Buffered data (>= 1 byte) exceeds the threshold
  // 3. CanWriteNewData() returns false
  std::string write_data = "foo";
  auto write_buf = base::MakeRefCounted<StringIOBuffer>(write_data);
  int rv = callback.adapter()->Write(write_buf.get(), write_buf->size(),
                                     callback.callback(),
                                     TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Resume mock data. The socket write fails with ERR_FAILED, causing
  // QUIC_PACKET_WRITE_ERROR which closes the connection and invokes OnClose().
  mock_quic_data_.Resume();
  rv = callback.WaitForResult();

  // Socket write error → QUIC_PACKET_WRITE_ERROR → MapQuicErrorToNetError()
  // returns ERR_QUIC_PROTOCOL_ERROR for any non-zero connection_error().
  EXPECT_THAT(rv, IsError(ERR_QUIC_PROTOCOL_ERROR));

  // Wait for the QUIC connection to fully close.
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));

  EXPECT_TRUE(mock_quic_data_.AllReadDataConsumed());
  EXPECT_TRUE(mock_quic_data_.AllWriteDataConsumed());
}

// Tests that the adapter correctly handles being destroyed from within its own
// write callback. This is a regression test to ensure OnClose() doesn't crash
// when the write callback destroys the adapter (use-after-free prevention).
TEST_P(WebSocketQuicStreamAdapterTest, WriteCallbackDestroysAdapter) {
  // Set a very low buffer threshold. When combined with flow control blocking,
  // any buffered data will exceed this threshold and cause Write() to return
  // ERR_IO_PENDING.
  quic::test::QuicFlagSaver flag_saver;
  SetQuicheFlag(quic_buffered_data_threshold, 1);

  int packet_number = 1;

  // Client sends SETTINGS during session initialization.
  mock_quic_data_.AddWrite(SYNCHRONOUS,
                           ConstructSettingsPacket(packet_number++));

  // Client sends WebSocket upgrade request headers.
  SpdyTestUtil spdy_util;
  quiche::HttpHeaderBlock request_header_block = WebSocketHttp2Request(
      "/", "www.example.org:443", "http://www.example.org", {});
  mock_quic_data_.AddWrite(
      SYNCHRONOUS,
      client_maker_.MakeRequestHeadersPacket(
          packet_number++, client_data_stream_id1_,
          /*fin=*/false, ConvertRequestPriorityToQuicPriority(LOWEST),
          std::move(request_header_block), nullptr));

  // Server sends response headers accepting the WebSocket upgrade.
  quiche::HttpHeaderBlock response_header_block = WebSocketHttp2Response({});
  mock_quic_data_.AddRead(
      ASYNC, server_maker_.MakeResponseHeadersPacket(
                 /*packet_number=*/1, client_data_stream_id1_, /*fin=*/false,
                 std::move(response_header_block),
                 /*spdy_headers_frame_length=*/nullptr));

  // Pause here to allow test to block flow control and issue a Write() before
  // processing further mock data.
  mock_quic_data_.AddRead(ASYNC, ERR_IO_PENDING);

  // When flow control is blocked and we try to write, QUIC sends a BLOCKED
  // frame to notify the peer. We need to expect this write.
  // Must be ASYNC because it follows a read pause.
  mock_quic_data_.AddWrite(
      ASYNC, client_maker_.Packet(packet_number++)
                 .AddAckFrame(1, 1, 1)
                 .AddFrame(quic::QuicFrame(
                     quic::QuicBlockedFrame(1, client_data_stream_id1_, 170)))
                 .Build());

  // Server closes the connection unexpectedly (simulates network error or
  // server shutdown). This triggers OnClose() which invokes the pending
  // write callback, which destroys the adapter via DeleterCallback.
  mock_quic_data_.AddRead(ASYNC, ERR_CONNECTION_CLOSED);

  EXPECT_CALL(mock_delegate_, OnHeadersReceived(_));

  Initialize();

  net::QuicChromiumClientSession::Handle* session_handle =
      GetQuicSessionHandle();
  ASSERT_TRUE(session_handle);

  // Create the WebSocket adapter.
  TestWebSocketQuicStreamAdapterCompletionCallback creation_callback;
  auto adapter = session_handle->CreateWebSocketQuicStreamAdapter(
      &mock_delegate_, creation_callback.callback(),
      TRAFFIC_ANNOTATION_FOR_TESTS);
  ASSERT_TRUE(adapter);
  EXPECT_TRUE(adapter->is_initialized());

  // Send WebSocket request headers to initiate the handshake.
  adapter->WriteHeaders(RequestHeaders(), false);

  // Start reading from the session and wait until the mock data is paused
  // (after receiving response headers but before connection close).
  session_->StartReading();
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return mock_quic_data_.GetSequencedSocketData()->IsPaused(); }));

  // Block flow control to force data to be buffered. Combined with the low
  // buffer threshold (1 byte), this causes CanWriteNewData() to return false
  // and Write() to return ERR_IO_PENDING.
  quic::QuicStream* quic_stream =
      session_->GetActiveStream(client_data_stream_id1_);
  ASSERT_TRUE(quic_stream);
  quic::QuicStreamOffset current_offset =
      quic::test::QuicStreamPeer::SendWindowOffset(quic_stream) -
      quic::test::QuicStreamPeer::SendWindowSize(quic_stream);
  quic::test::QuicStreamPeer::SetSendWindowOffset(quic_stream, current_offset);

  // DeleterCallback will destroy the adapter when the write completes.
  DeleterCallback callback(std::move(adapter));

  // Issue an async write. The callback will destroy the adapter when invoked.
  // Write returns ERR_IO_PENDING because:
  // 1. Flow control is blocked, so data is buffered (not sent)
  // 2. Buffered data (>= 1 byte) exceeds the threshold
  // 3. CanWriteNewData() returns false
  std::string write_data = "foo";
  auto write_buf = base::MakeRefCounted<StringIOBuffer>(write_data);
  int rv = callback.adapter()->Write(write_buf.get(), write_buf->size(),
                                     callback.callback(),
                                     TRAFFIC_ANNOTATION_FOR_TESTS);
  EXPECT_THAT(rv, IsError(ERR_IO_PENDING));

  // Resume the mock data to deliver the connection close event.
  // This triggers the write callback which destroys the adapter.
  mock_quic_data_.Resume();
  rv = callback.WaitForResult();
  EXPECT_THAT(rv, IsError(ERR_QUIC_PROTOCOL_ERROR));

  // Wait for the QUIC connection to fully close.
  ASSERT_TRUE(base::test::RunUntil(
      [&]() { return !session_->connection()->connected(); }));

  EXPECT_TRUE(mock_quic_data_.AllReadDataConsumed());
  EXPECT_TRUE(mock_quic_data_.AllWriteDataConsumed());
}

}  // namespace net::test
