diff --git a/CMakeLists.txt b/CMakeLists.txt index b6840b34a4ab..ede7424c0765 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1439,17 +1439,6 @@ if(BUILD_TESTS) ADDITIONAL_ARGS --js-app-bundle ${CMAKE_SOURCE_DIR}/samples/apps/logging/js ) - add_e2e_test( - NAME e2e_logging_http2 - PYTHON_SCRIPT ${CMAKE_SOURCE_DIR}/tests/e2e_logging.py - BUCKET bucket_b - BUILD_DEPENDS logging logging_cose_only programmability js_generic - ADDITIONAL_ARGS - --js-app-bundle - ${CMAKE_SOURCE_DIR}/samples/apps/logging/js - --http2 - ) - add_e2e_test( NAME partitions PYTHON_SCRIPT ${CMAKE_SOURCE_DIR}/tests/partitions_test.py diff --git a/doc/use_apps/issue_commands.rst b/doc/use_apps/issue_commands.rst index a3adb4de7005..32e7939a4c37 100644 --- a/doc/use_apps/issue_commands.rst +++ b/doc/use_apps/issue_commands.rst @@ -43,7 +43,7 @@ require setting ``ccf.gov.msg.type``, ``ccf.gov.msg.created_at``, and optionally A signing script (``ccf_cose_sign1``) is provided as part of the `ccf Python package `_. The output can be piped directly into curl, or any other HTTP client. -Commands can also be signed using the `python-cwt `_ library, and sent with any standard HTTP library such as `Python HTTPX `_. +Commands can also be signed using the `python-cwt `_ library, and sent with any standard HTTP library such as `Python Requests `_. Idempotence ^^^^^^^^^^^ diff --git a/scripts/fetch_amd_collateral.py b/scripts/fetch_amd_collateral.py index b2d63e01ac98..cfd5c442c329 100644 --- a/scripts/fetch_amd_collateral.py +++ b/scripts/fetch_amd_collateral.py @@ -5,7 +5,7 @@ from enum import Enum import logging import sys -import httpx +import requests import base64 from cryptography import x509 from cryptography.hazmat.primitives import serialization @@ -131,10 +131,8 @@ def make_chain_url(base_url, product_family): ) logging.info(f"Fetching AMD leaf cert from {leaf_url}") - with httpx.Client() as client: - leaf_response = client.get( - leaf_url, - ) + with requests.Session() as client: + leaf_response = client.get(leaf_url, timeout=30) leaf_response.raise_for_status() der = leaf_response.content leaf = ( @@ -147,8 +145,8 @@ def make_chain_url(base_url, product_family): chain_url = make_chain_url(args.base_url, args.product_family) logging.info(f"Fetching AMD chain cert from {chain_url}") - with httpx.Client() as client: - chain_response = client.get(chain_url) + with requests.Session() as client: + chain_response = client.get(chain_url, timeout=30) chain_response.raise_for_status() chain = chain_response.text logging.info(f"AMD chain cert response: {chain_response.text}") diff --git a/scripts/requirements.txt b/scripts/requirements.txt index 6cfcc7a76cf9..6cb3c99fbd84 100644 --- a/scripts/requirements.txt +++ b/scripts/requirements.txt @@ -1,3 +1,3 @@ cryptography >= 48.0.1, < 49 -httpx >= 0.28.1, < 0.29 +requests >= 2.32.5, < 3 GitPython >= 3.1.45, < 4 diff --git a/tests/ci-buckets.txt b/tests/ci-buckets.txt index 8b363da15f70..d460904c2acd 100644 --- a/tests/ci-buckets.txt +++ b/tests/ci-buckets.txt @@ -8,7 +8,6 @@ bucket_b: recovery_stale_snapshot_join_test recovery_intermediate_snapshot_join_test recovery_snapshot_endorsements_test - e2e_logging_http2 schema_test nodes_test diff --git a/tests/client_protocols.py b/tests/client_protocols.py index 10ada24ad7bb..1a3c17a97a3e 100644 --- a/tests/client_protocols.py +++ b/tests/client_protocols.py @@ -6,6 +6,7 @@ import subprocess import infra.e2e_args +import infra.interfaces import infra.net import infra.network import infra.proc @@ -14,6 +15,10 @@ # As installed by setup scripts H2SPEC_BIN = "/opt/h2spec/h2spec" +# HTTP/2 is only served on a dedicated interface: the Python test client +# only speaks HTTP/1.1, which the primary interface keeps for governance +HTTP2_RPC_INTERFACE = "http2_interface" + def compare_golden(): script_path = os.path.realpath(__file__) @@ -114,9 +119,9 @@ def test_http2(network, args): "--tls", "--insecure", "--host", - node.get_public_rpc_host(), + node.get_public_rpc_host(HTTP2_RPC_INTERFACE), "--port", - f"{node.get_public_rpc_port()}", + f"{node.get_public_rpc_port(HTTP2_RPC_INTERFACE)}", "--strict", ], check=True, @@ -124,17 +129,35 @@ def test_http2(network, args): assert r.returncode == 0 +def single_interface_node(args): + # Retain only the primary interface, delete any others + nodes = infra.e2e_args.nodes(args, 1) + nodes[0].rpc_interfaces = { + infra.interfaces.PRIMARY_RPC_INTERFACE: nodes[0].get_primary_interface() + } + return nodes + + def run(args): + # The TLS report is generated against a node with a single (HTTP/1.1) + # interface, and should still mention ALPN HTTP/1.1 as HTTP/2 is + # experimental as of 3.x + args.nodes = single_interface_node(args) with infra.network.network( args.nodes, args.binary_dir, args.debug_nodes, pdb=args.pdb ) as network: network.start_and_open(args) test_tls(network, args) - # Note: Start new network with HTTP/2 as TLS report should still - # mention ALPN HTTP/1.1 as HTTP/2 is experimental as of 3.x - args.http2 = True - args.nodes = infra.e2e_args.nodes(args, 1) + # Start a new network with an additional, dedicated HTTP/2 interface for + # the compliance test. The primary interface stays HTTP/1.1, so that the + # service can be opened with the Python client. + args.nodes = single_interface_node(args) + primary_interface = args.nodes[0].get_primary_interface() + args.nodes[0].rpc_interfaces[HTTP2_RPC_INTERFACE] = infra.interfaces.RPCInterface( + host=primary_interface.host, + app_protocol="HTTP2", + ) with infra.network.network( args.nodes, args.binary_dir, args.debug_nodes, pdb=args.pdb ) as network: @@ -146,11 +169,4 @@ def run(args): args = infra.e2e_args.cli_args() args.package = "samples/apps/logging/logging" - args.nodes = infra.e2e_args.nodes(args, 1) - - # Retain only the primary interface, delete any others - args.nodes[0].rpc_interfaces = { - infra.interfaces.PRIMARY_RPC_INTERFACE: args.nodes[0].get_primary_interface() - } - run(args) diff --git a/tests/connections.py b/tests/connections.py index 944cab7f70c4..806ce69e2a47 100644 --- a/tests/connections.py +++ b/tests/connections.py @@ -12,7 +12,6 @@ import time import fuzzing -import httpx import infra.checker import infra.e2e_args import infra.interfaces @@ -217,11 +216,6 @@ def create_connections_until_exhaustion( client_fn( identity="user0", connection_timeout=1, - limits=httpx.Limits( - max_connections=1, - max_keepalive_connections=1, - keepalive_expiry=30, - ), ) ) ) @@ -269,10 +263,20 @@ def create_connections_until_exhaustion( LOG.success( f"{primary_pid} has {num_fds}/{max_fds} open file descriptors" ) - r = clients[0].get("/node/metrics") - assert r.status_code == http.HTTPStatus.OK, r.status_code - peak_metrics = r.body.json()["sessions"] - assert peak_metrics["active"] <= peak_metrics["peak"], peak_metrics + # Sessions refused with a 503 are closed asynchronously by the + # node, so briefly wait for the active count to settle + end_time = time.time() + 3 + while True: + r = clients[0].get("/node/metrics") + assert r.status_code == http.HTTPStatus.OK, r.status_code + peak_metrics = r.body.json()["sessions"] + assert peak_metrics["active"] <= peak_metrics["peak"], peak_metrics + if ( + peak_metrics["active"] == len(healthy_clients) + or time.time() > end_time + ): + break + time.sleep(0.1) assert peak_metrics["active"] == len(healthy_clients), ( peak_metrics, len(healthy_clients), @@ -352,6 +356,11 @@ def create_connections_until_exhaustion( resource.prlimit(primary_pid, resource.RLIMIT_NOFILE, (max_fds, max_fds)) LOG.success(f"Setting max fds to dangerously low {max_fds} on {primary_pid}") + # The node is expected to crash when it runs out of file descriptors. + # That may only happen after the client has seen its responses, as + # ledger writes are asynchronous, so tolerate fatal errors at shutdown + # regardless of whether the crash is observed by the client. + network.ignore_errors_on_shutdown() try: num_fds = create_connections_until_exhaustion(to_create) except Exception as e: @@ -359,7 +368,6 @@ def create_connections_until_exhaustion( f"Node with only {max_fds} fds crashed when allowed to created {args.max_open_sessions} sessions, as expected" ) LOG.warning(e) - network.ignore_errors_on_shutdown() else: LOG.warning("Expected a fatal crash and saw none!") diff --git a/tests/e2e_common_endpoints.py b/tests/e2e_common_endpoints.py index a7454a98a5d4..3c74fa05a249 100644 --- a/tests/e2e_common_endpoints.py +++ b/tests/e2e_common_endpoints.py @@ -248,16 +248,12 @@ def run_large_message_test( expected_errors[metrics_name] += 1 assert get_main_interface_errors() == expected_errors - def get_sizes(n, http2): - ns = [n // 2, n - 10, n - 1, n, n + 1, n + 10, n * 2] - if not http2: - # nghttp2 does not currently allow header larger than 64KB - # https://github.com/nghttp2/nghttp2/issues/1841 - ns.append(n * 20) + def get_sizes(n): + ns = [n // 2, n - 10, n - 1, n, n + 1, n + 10, n * 2, n * 20] random.shuffle(ns) return ns - for s in get_sizes(args.max_http_body_size, args.http2): + for s in get_sizes(args.max_http_body_size): long_msg = "X" * s LOG.info(f"Verifying cap on max body size, sending a {s} byte body") run_large_message_test( @@ -270,7 +266,7 @@ def get_sizes(n, http2): headers={"content-type": "application/json"}, ) - for s in get_sizes(args.max_http_header_size, args.http2): + for s in get_sizes(args.max_http_header_size): long_header = "X" * s LOG.info(f"Verifying cap on max header value, sending a {s} byte header value") run_large_message_test( @@ -292,36 +288,35 @@ def get_sizes(n, http2): headers={long_header: "some header value"}, ) - if not args.http2: - for size in ( - args.max_http_request_target_size - 1, - args.max_http_request_target_size, - args.max_http_request_target_size + 1, - ): - prefix = "/node/commit?padding=" - if size < len(prefix): - LOG.warning( - f"Skipping {size} byte request target: the test endpoint " - f"requires at least {len(prefix)} bytes" - ) - continue - target = prefix + "a" * (size - len(prefix)) - assert len(target) == size - LOG.info(f"Verifying cap on request target, sending a {size} byte target") - run_large_message_test( - args.max_http_request_target_size, - http.HTTPStatus.REQUEST_URI_TOO_LONG, - "RequestTargetTooLong", - "request_target_too_long", - len(target), - path=target, + for size in ( + args.max_http_request_target_size - 1, + args.max_http_request_target_size, + args.max_http_request_target_size + 1, + ): + prefix = "/node/commit?padding=" + if size < len(prefix): + LOG.warning( + f"Skipping {size} byte request target: the test endpoint " + f"requires at least {len(prefix)} bytes" ) + continue + target = prefix + "a" * (size - len(prefix)) + assert len(target) == size + LOG.info(f"Verifying cap on request target, sending a {size} byte target") + run_large_message_test( + args.max_http_request_target_size, + http.HTTPStatus.REQUEST_URI_TOO_LONG, + "RequestTargetTooLong", + "request_target_too_long", + len(target), + path=target, + ) # Note: infra generally inserts extra headers (eg, content type and length, user-agent, accept) - extra_headers_count = infra.clients.CCFClient.default_impl_type.extra_headers_count( - args.http2 + extra_headers_count = ( + infra.clients.CCFClient.default_impl_type.extra_headers_count() ) - for s in get_sizes(args.max_http_headers_count, args.http2): + for s in get_sizes(args.max_http_headers_count): LOG.info(f"Verifying on cap on max headers count, sending {s} headers") headers = {f"header-{h}": str(h) for h in range(s - extra_headers_count)} run_large_message_test( diff --git a/tests/e2e_logging.py b/tests/e2e_logging.py index a6ef251cc6a7..949b7f94a374 100644 --- a/tests/e2e_logging.py +++ b/tests/e2e_logging.py @@ -40,7 +40,7 @@ from cryptography.hazmat.backends import default_backend from cryptography.x509 import ObjectIdentifier, load_pem_x509_certificate from infra.log_capture import flush_info -from infra.member import AckException, RecoveryRole +from infra.member import RecoveryRole from infra.runner import ConcurrentRunner from infra.tx_status import TxStatus from loguru import logger as LOG @@ -157,13 +157,11 @@ def test(network, args): network=network, number_txs=1, ) - # HTTP2 doesn't support forwarding - if not args.http2: - network.txs.issue( - network=network, - number_txs=1, - on_backup=True, - ) + network.txs.issue( + network=network, + number_txs=1, + on_backup=True, + ) network.txs.verify() return network @@ -205,14 +203,7 @@ def send_raw_content(content): def send_bad_raw_content(content): nonlocal additional_parsing_errors - try: - response = send_raw_content(content) - except http.client.RemoteDisconnected: - assert args.http2, "HTTP/2 interface should close session without error" - additional_parsing_errors += 1 - return - else: - assert not args.http2, "HTTP/1.1 interface should return valid error" + response = send_raw_content(content) response_body = response.read() LOG.warning(response_body) @@ -260,26 +251,22 @@ def send_corrupt_variations(content): == initial_parsing_errors + additional_parsing_errors ) - if not args.http2: - good_content = b"GET /node/state HTTP/1.1\r\n\r\n" - response = send_raw_content(good_content) - assert response.status == http.HTTPStatus.OK, (response.status, response.read()) - send_corrupt_variations(good_content) + good_content = b"GET /node/state HTTP/1.1\r\n\r\n" + response = send_raw_content(good_content) + assert response.status == http.HTTPStatus.OK, (response.status, response.read()) + send_corrupt_variations(good_content) # Valid transactions are still accepted network.txs.issue( network=network, number_txs=1, ) - - # HTTP/2 does not support forwarding - if not args.http2: - network.txs.issue( - network=network, - number_txs=1, - on_backup=True, - ) - network.txs.verify() + network.txs.issue( + network=network, + number_txs=1, + on_backup=True, + ) + network.txs.verify() return network @@ -346,7 +333,7 @@ def parse_result_out(r): ) expected_response_body, status_code, http_version = parse_result_out(res) assert status_code == "200", status_code - assert http_version == "2" if args.http2 else "1.1", http_version + assert http_version == "1.1", http_version protocols = { # WebSockets upgrade request is ignored @@ -366,30 +353,14 @@ def parse_result_out(r): "option --http3: is unknown", ] }, + # HTTP/1.x requests succeed, as HTTP/1.1 + "--http1.0": {}, + "--http1.1": {}, + # TLS handshake negotiates HTTP/1.1 + "--http2": {}, + # This is disabled because the behaviour of curl differs from version 8.10, so we do not get consistent results across platforms + # "--http2-prior-knowledge": {}, } - if args.http2: - protocols.update( - { - # HTTP/1.x requests fail with closed connection, as HTTP/2 - "--http1.0": {"errors": ["Empty reply from server"]}, - "--http1.1": {"errors": ["Empty reply from server"]}, - # TLS handshake negotiates HTTP/2 - "--http2": {}, - "--http2-prior-knowledge": {}, - } - ) - else: # HTTP/1.1 - protocols.update( - { - # HTTP/1.x requests succeed, as HTTP/1.1 - "--http1.0": {}, - "--http1.1": {}, - # TLS handshake negotiates HTTP/1.1 - "--http2": {}, - # This is disabled because the behaviour of curl differs from version 8.10, so we do not get consistent results across platforms - # "--http2-prior-knowledge": {}, - } - ) # Test additional protocols with curl for protocol, expected_result in protocols.items(): @@ -406,7 +377,7 @@ def parse_result_out(r): response_body == expected_response_body ), f"{response_body}\n !=\n{expected_response_body}" assert status_code == "200", status_code - assert http_version == "2" if args.http2 else "1.1", http_version + assert http_version == "1.1", http_version else: assert res.returncode != 0, res.returncode err = res.stderr.decode() @@ -418,14 +389,12 @@ def parse_result_out(r): network=network, number_txs=1, ) - # HTTP/2 does not support forwarding - if not args.http2: - network.txs.issue( - network=network, - number_txs=1, - on_backup=True, - ) - network.txs.verify() + network.txs.issue( + network=network, + number_txs=1, + on_backup=True, + ) + network.txs.verify() return network @@ -709,13 +678,7 @@ def require_new_response(r): def test_custom_auth(network, args): primary, other = network.find_primary_and_any_backup() - nodes = (primary, other) - - if args.http2: - # HTTP2 doesn't support forwarding - nodes = (primary,) - - for node in nodes: + for node in (primary, other): with node.client() as c: LOG.info("Request without custom headers is refused") r = c.get("/app/custom_auth") @@ -755,13 +718,7 @@ def test_custom_auth(network, args): def test_custom_auth_safety(network, args): primary, other = network.find_primary_and_any_backup() - nodes = (primary, other) - - if args.http2: - # HTTP2 doesn't support forwarding - nodes = (primary,) - - for node in nodes: + for node in (primary, other): with node.client() as c: r = c.get( "/app/custom_auth", @@ -1544,50 +1501,27 @@ def escaped_query_tests(c, endpoint): @reqs.description("Testing forwarding on member and user frontends") @reqs.supports_methods("/app/log/private") @reqs.at_least_n_nodes(2) -@reqs.no_http2() @app.scoped_txs() def test_forwarding_frontends(network, args): backup = network.find_any_backup() - try: - with backup.client() as c: - check_commit = infra.checker.Checker(c) - ack = network.consortium.get_any_active_member().ack(backup) - check_commit(ack) - except AckException as e: - assert args.http2 is True - assert e.response.status_code == http.HTTPStatus.NOT_IMPLEMENTED - r = e.response.body.json() - assert ( - r["error"]["message"] - == "Request cannot be forwarded to primary on HTTP/2 interface." - ), r - else: - assert args.http2 is False + with backup.client() as c: + check_commit = infra.checker.Checker(c) + ack = network.consortium.get_any_active_member().ack(backup) + check_commit(ack) - try: - msg = "forwarded_msg" - log_id = 7 - network.txs.issue( - network, - number_txs=1, - on_backup=True, - idx=log_id, - send_public=False, - msg=msg, - ) - except infra.logging_app.LoggingTxsIssueException as e: - assert args.http2 is True - assert e.response.status_code == http.HTTPStatus.NOT_IMPLEMENTED - r = e.response.body.json() - assert ( - r["error"]["message"] - == "Request cannot be forwarded to primary on HTTP/2 interface." - ), r - else: - assert args.http2 is False + msg = "forwarded_msg" + log_id = 7 + network.txs.issue( + network, + number_txs=1, + on_backup=True, + idx=log_id, + send_public=False, + msg=msg, + ) - if args.package.startswith("samples/apps/logging/logging") and not args.http2: + if args.package.startswith("samples/apps/logging/logging"): with backup.client("user0") as c: escaped_query_tests(c, "request_query") @@ -1596,7 +1530,6 @@ def test_forwarding_frontends(network, args): @reqs.description("Testing forwarding on user frontends without actor app prefix") @reqs.at_least_n_nodes(2) -@reqs.no_http2() def test_forwarding_frontends_without_app_prefix(network, args): msg = "forwarded_msg" log_id = 7 @@ -1616,7 +1549,6 @@ def test_forwarding_frontends_without_app_prefix(network, args): @reqs.description("Testing forwarding on long-lived connection") @reqs.supports_methods("/app/log/private") @reqs.at_least_n_nodes(2) -@reqs.no_http2() def test_long_lived_forwarding(network, args): primary, _ = network.find_primary() @@ -1624,15 +1556,19 @@ def test_long_lived_forwarding(network, args): new_node = network.create_node() # Message limit must be high enough that the hard limit will not be reached - # by the combined work of all threads. Note that each thread produces multiple - # node-to-node messages - a forwarded write and response, Raft AEs. If these - # arrive too fast, they will trigger the hard cap and the node-to-node keys - # will be reset, potentially invalidating in-flight messages and causing client + # by the combined work of all threads between the point where the soft limit + # (half the message limit) triggers a key exchange and the point where that + # exchange completes. Note that each thread produces multiple node-to-node + # messages - a forwarded write and response, Raft AEs. If these arrive too + # fast, they will trigger the hard cap and the node-to-node keys will be + # reset, potentially invalidating in-flight messages and causing client # requests to time out. This margin depends on client request rate, so must # stay comfortably above the concurrent in-flight message burst produced by # n_threads clients sending as fast as the network allows. n_threads = 5 - message_limit = 90 + message_limit = 400 + # Enough requests per thread for several key rotations to happen + requests_per_thread = 300 new_node_args = copy.deepcopy(args) new_node_args.node_to_node_message_limit = message_limit @@ -1664,7 +1600,7 @@ def fn(worker_id, request_count, should_log): threads.append( threading.Thread( target=fn, - args=(i, 3 * message_limit, i == 0), + args=(i, requests_per_thread, i == 0), name=f"{current_thread_name}:worker-{i}", ) ) @@ -2370,7 +2306,6 @@ def additional_interfaces(local_node_id): for interface_name, host in additional_interfaces(local_node_id).items(): node_host.rpc_interfaces[interface_name] = infra.interfaces.RPCInterface( host=host, - app_protocol="HTTP2" if args.http2 else "HTTP1", ) txs = app.LoggingTxs("user0") @@ -2654,12 +2589,10 @@ def do_main_tests(network, args): test_cose_signature_schema(network, args) test_cose_receipt_schema(network, args) - # HTTP2 doesn't support forwarding - if not args.http2: - test_forwarding_frontends(network, args) - test_forwarding_frontends_without_app_prefix(network, args) - if not os.getenv("TSAN_OPTIONS"): - test_long_lived_forwarding(network, args) + test_forwarding_frontends(network, args) + test_forwarding_frontends_without_app_prefix(network, args) + if not os.getenv("TSAN_OPTIONS"): + test_long_lived_forwarding(network, args) test_user_data_ACL(network, args) test_cert_prefix(network, args) test_anonymous_caller(network, args) @@ -2689,13 +2622,12 @@ def do_main_tests(network, args): if args.package.startswith("samples/apps/logging/logging"): test_etags(network, args) test_cose_config(network, args) - if not args.http2: - test_blocking_calls(network, args) + test_blocking_calls(network, args) # These tests require a service which has only ever emitted COSE signatures, # unlike a service upgraded from Dual mode with legacy signatures in its ledger. is_cose_only_from_genesis = args.package.endswith("_cose_only") - if is_cose_only_from_genesis and not args.http2: + if is_cose_only_from_genesis: test_cose_set_member(network, args) diff --git a/tests/infra/clients.py b/tests/infra/clients.py index 1f97d7fce7d5..9c5861442bef 100644 --- a/tests/infra/clients.py +++ b/tests/infra/clients.py @@ -24,8 +24,15 @@ from typing import Any import ccf.cose -import httpcore.backends.sync -import httpx +import requests +import requests.adapters +import requests.auth +import requests.exceptions +import urllib3 +import urllib3._collections +import urllib3.connection +import urllib3.connectionpool +import urllib3.exceptions from ccf.tx_id import TxID from cryptography import x509 from cryptography.hazmat.backends import default_backend @@ -78,9 +85,7 @@ def get_clock(): return _per_thread.CLOCK -class HttpSig(httpx.Auth): - requires_request_body = True - +class HttpSig(requests.auth.AuthBase): def __init__(self, key_id, pem_private_key): self.key_id = key_id self.private_key = load_pem_private_key( @@ -109,16 +114,19 @@ def add_signature_headers(headers, content, method, path, key_id, private_key): f'Signature keyId="{key_id}",algorithm="hs2019",headers="(request-target) digest content-length",signature="{b64signature}"' ) - def auth_flow(self, request): + def __call__(self, request: requests.PreparedRequest) -> requests.PreparedRequest: + content = request.body if request.body is not None else b"" + if isinstance(content, str): + content = content.encode("utf-8") HttpSig.add_signature_headers( request.headers, - request.content, + content, request.method, - request.url.raw_path.decode("utf-8"), + request.path_url, self.key_id, self.private_key, ) - yield request + return request def truncate(string: str, max_len: int = 256): @@ -213,17 +221,19 @@ def __repr__(self): class RequestsResponseBody(ResponseBody): - def __init__(self, response: httpx.Response): + def __init__(self, response: requests.Response): self._response = response def data(self): return self._response.content def text(self): - return self._response.text + # Decode as UTF-8 regardless of any charset guess made by requests, + # replacing invalid sequences rather than raising + return self._response.content.decode("utf-8", errors="replace") def json(self): - return self._response.json() + return json.loads(self._response.content) class RawResponseBody(ResponseBody): @@ -470,9 +480,6 @@ def __init__( else: self.ca_curve = None self.protocol = kwargs.get("protocol") if "protocol" in kwargs else "https" - self.extra_args = [] - if kwargs.get("http2"): - self.extra_args.append("--http2") self.cose_header_builder = cose_protected_headers_api_classic def request( @@ -564,9 +571,6 @@ def request( if not self.ca and not self.session_auth: cmd.extend(["-k"]) # Allow insecure connections - for arg in self.extra_args: - cmd.append(arg) - cmd_s = " ".join(cmd) env = {k: v for k, v in os.environ.items()} @@ -608,41 +612,62 @@ def close(self): pass @staticmethod - def extra_headers_count(http2=False): + def extra_headers_count(): # curl inserts the following headers in every request - if http2: - # :method: GET/POST - # :authority:
- # :scheme: https - # :path: /path - # accept: */* - # user-agent: curl/ - return 6 - else: - # host:
- # user-agent: curl/ - # accept: */* - return 3 + # host:
+ # user-agent: curl/ + # accept: */* + return 3 -class _NoDelaySyncBackend(httpcore.backends.sync.SyncBackend): +class ConnectionEstablishmentError(urllib3.exceptions.NewConnectionError): """ - Sync network backend that enables TCP_NODELAY, avoiding a ~40ms - Nagle/delayed-ACK stall per request on the pinned httpcore 0.16 - (newer httpcore versions set this by default). + Raised for any failure while establishing a connection (TCP connect or TLS + handshake), so that it can be told apart from I/O errors on a connection + which was already established. """ - def connect_tcp(self, *args, **kwargs): - stream = super().connect_tcp(*args, **kwargs) - sock = stream.get_extra_info("socket") - if sock is not None: - sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) - return stream + +class _ClassifyingHTTPSConnection(urllib3.connection.HTTPSConnection): + def connect(self): + try: + super().connect() + except ( + urllib3.exceptions.NewConnectionError, + urllib3.exceptions.ConnectTimeoutError, + TimeoutError, + ): + raise + except Exception as exc: + raise ConnectionEstablishmentError( + self, f"Failed to establish a new connection: {exc}" + ) from exc + + +class _ClassifyingHTTPSConnectionPool(urllib3.connectionpool.HTTPSConnectionPool): + ConnectionCls = _ClassifyingHTTPSConnection + + +class _CCFHTTPAdapter(requests.adapters.HTTPAdapter): + def init_poolmanager(self, connections, maxsize, block=False, **pool_kwargs): + super().init_poolmanager(connections, maxsize, block=block, **pool_kwargs) + self.poolmanager.pool_classes_by_scheme = { + **self.poolmanager.pool_classes_by_scheme, + "https": _ClassifyingHTTPSConnectionPool, + } + # Close pooled connections as soon as the pool manager is cleared (ie + # when the session is closed). urllib3 >= 2.8 instead waits for the + # pools to be garbage collected, but responses keep a reference to + # their pool, so a retained response would keep its connection (and + # the corresponding CCF session) open after the client is closed. + self.poolmanager.pools = urllib3._collections.RecentlyUsedContainer( + connections, dispose_func=lambda pool: pool.close() + ) -class HttpxClient: +class RequestsClient: """ - CCF default client and wrapper around Python httpx, handling HTTP signatures. + CCF default client and wrapper around Python requests, handling HTTP signatures. """ _auth_provider = HttpSig @@ -653,13 +678,19 @@ class HttpxClient: def __init__( self, hostname: str, - ca: str, + ca: str | None, session_auth: Identity | None = None, signing_auth: Identity | None = None, cose_signing_auth: Identity | None = None, common_headers: dict | None = None, + protocol: str = "https", + headers: dict | None = None, **kwargs, ): + if kwargs: + raise TypeError( + f"Unexpected RequestsClient arguments: {', '.join(sorted(kwargs))}" + ) self.hostname = hostname self.ca = ca self.session_auth = session_auth @@ -667,24 +698,20 @@ def __init__( self.cose_signing_auth = cose_signing_auth self.common_headers = common_headers self.key_id = None - cert = None + self.protocol = protocol + self.session = requests.Session() + adapter = _CCFHTTPAdapter() + self.session.mount("https://", adapter) + self.session.mount("http://", adapter) + if self.ca: + self.session.verify = self.ca + else: + self.session.verify = False + urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) if self.session_auth: - cert = (self.session_auth.cert, self.session_auth.key) - self.protocol = "https" - if "protocol" in kwargs: - self.protocol = kwargs.get("protocol") - kwargs.pop("protocol") - self.session = httpx.Client(verify=self.ca, cert=cert, **kwargs) - # Swap in a network backend which sets TCP_NODELAY (see - # _NoDelaySyncBackend), regardless of whether the transport was - # constructed for HTTP/1.1 or HTTP/2. - pool = getattr( - self.session._transport, "_pool", None - ) # pylint: disable=protected-access - if pool is not None: - pool._network_backend = ( - _NoDelaySyncBackend() - ) # pylint: disable=protected-access + self.session.cert = (self.session_auth.cert, self.session_auth.key) + if headers: + self.session.headers.update(headers) sig_auth = signing_auth or cose_signing_auth if sig_auth: with open(sig_auth.cert, encoding="utf-8") as cert_file: @@ -700,30 +727,50 @@ def __init__( def _request( self, request: Request, - request_body: bytes, + request_body: bytes | None, auth: Any, extra_headers: dict | None, timeout: int, ): self._last_request = (request, request_body, auth, extra_headers, timeout) try: + # Redirects are followed by CCFClient, so that the request is + # re-issued unmodified (including auth) against the new target response = self.session.request( request.http_verb, url=f"{self.protocol}://{self.hostname}{request.path}", auth=auth, headers=extra_headers, timeout=timeout, - content=request_body, + data=request_body, + allow_redirects=False, ) - except httpx.TimeoutException as exc: + except requests.exceptions.Timeout as exc: raise TimeoutError from exc - except httpx.ConnectError as exc: - raise CCFConnectionException from exc - except (httpx.WriteError, httpx.ReadError, httpx.RemoteProtocolError) as exc: + except requests.exceptions.ChunkedEncodingError as exc: + raise CCFIOException from exc + except requests.exceptions.ConnectionError as exc: + # requests wraps both connection establishment failures and I/O + # errors on an established connection in ConnectionError, so + # inspect the underlying urllib3 error to tell them apart. Only + # the former should be retried by CCFClient. + cause = exc.args[0] if exc.args else None + if isinstance(cause, urllib3.exceptions.MaxRetryError): + cause = cause.reason + if isinstance(cause, urllib3.exceptions.ReadTimeoutError): + raise TimeoutError from exc + if isinstance( + cause, + ( + urllib3.exceptions.NewConnectionError, + urllib3.exceptions.ConnectTimeoutError, + ), + ): + raise CCFConnectionException from exc raise CCFIOException from exc except Exception as exc: raise RuntimeError( - f"HttpxClient failed with unexpected error: {exc}" + f"RequestsClient failed with unexpected error: {exc}" ) from exc return Response.from_requests_response(response) @@ -779,7 +826,10 @@ def request( request_body = json.dumps(request.body).encode() content_type = CONTENT_TYPE_JSON - if "content-type" not in request.headers and len(request.body) > 0: + # requests treats header names case-insensitively, so a caller's + # Content-Type must be detected regardless of its case + has_content_type = any(k.lower() == "content-type" for k in extra_headers) + if not has_content_type and len(request.body) > 0: extra_headers["content-type"] = content_type if self.cose_signing_auth is not None and request.http_verb != "GET": @@ -830,24 +880,14 @@ def close(self): self.session.close() @staticmethod - def extra_headers_count(http2=False): - # httpx inserts the following headers in every request - if http2: - # :method: GET/POST - # :authority:
- # :scheme: https - # :path: /path - # accept: */* - # accept-encoding: gzip, deflate, br - # user-agent: python-httpx/ - return 7 - else: - # host:
- # accept: */* - # accept-encoding: gzip, deflate, br - # connection: keep-alive - # user-agent: python-httpx/ - return 5 + def extra_headers_count(): + # requests inserts the following headers in every request + # host:
+ # accept: */* + # accept-encoding: gzip, deflate, ... + # connection: keep-alive + # user-agent: python-requests/ + return 5 class RawSocketClient: @@ -1040,13 +1080,13 @@ class CCFClient: default_impl_type = ( CurlClient if os.getenv("CURL_CLIENT") - else RawSocketClient if os.getenv("SOCKET_CLIENT") else HttpxClient + else RawSocketClient if os.getenv("SOCKET_CLIENT") else RequestsClient ) def set_created_at_override(self, value): if isinstance(self.client_impl, CurlClient): assert value.tzinfo, "created_at must be timezone aware" - elif isinstance(self.client_impl, HttpxClient): + elif isinstance(self.client_impl, RequestsClient): assert ( isinstance(value, int) or value.tzinfo ), "created_at must be integer or timezone aware" @@ -1063,7 +1103,7 @@ def __init__( connection_timeout: int = DEFAULT_CONNECTION_TIMEOUT_SEC, election_timeout_ms: int | None = None, description: str | None = None, - impl_type: CurlClient | HttpxClient | RawSocketClient = default_impl_type, + impl_type: CurlClient | RequestsClient | RawSocketClient = default_impl_type, common_headers: dict | None = None, openapi_validator=None, **kwargs, @@ -1149,9 +1189,14 @@ def _call( r = Request(redirect_path, body, http_verb, headers) request_client = temp_client - response = request_client.request( - r, timeout, cose_header_parameters_override - ) + # Response bodies are fully read by the client implementation, so + # the temporary client (and its connection) can be closed at once + try: + response = request_client.request( + r, timeout, cose_header_parameters_override + ) + finally: + temp_client.close() flush_info([str(response)], log_capture, 3) if self.openapi_validator is not None and validate_openapi: diff --git a/tests/infra/e2e_args.py b/tests/infra/e2e_args.py index 729bbe5197d0..2c9e0f2af577 100644 --- a/tests/infra/e2e_args.py +++ b/tests/infra/e2e_args.py @@ -83,7 +83,6 @@ "max_http_headers_count": ( "network.rpc_interfaces.*.http_configuration.max_headers_count" ), - "http2": "network.rpc_interfaces.*.app_protocol", "snp_endorsements_servers": None, "forwarding_timeout_ms": "network.rpc_interfaces.*.forwarding_timeout_ms", "tick_ms": "tick_interval", @@ -147,7 +146,6 @@ def _convert_curve_id(value): "max_http_body_size": _convert_size_string_to_bytes, "max_http_header_size": _convert_size_string_to_bytes, "max_http_request_target_size": _convert_size_string_to_bytes, - "http2": lambda value: value == "HTTP2", "tick_ms": lambda value: _convert_time_string(value, "ms"), } @@ -658,12 +656,6 @@ def cli_args( default=256, type=int, ) - parser.add_argument( - "--http2", - help="Enable HTTP/2 for all interfaces", - action="store_true", - default=False, - ) parser.add_argument( "--snp-endorsements-servers", help="Servers used to retrieve attestation report endorsement certificates (AMD SEV-SNP only)", diff --git a/tests/infra/interfaces.py b/tests/infra/interfaces.py index aa0477b1193c..746583d0362a 100644 --- a/tests/infra/interfaces.py +++ b/tests/infra/interfaces.py @@ -211,7 +211,6 @@ def apply_args(self, args): self.max_http_request_target_size = args.max_http_request_target_size self.max_http_headers_count = args.max_http_headers_count self.forwarding_timeout_ms = args.forwarding_timeout_ms - self.app_protocol = "HTTP2" if args.http2 else "HTTP1" def parse_from_str(self, s): # Format: local|ssh(,tcp|udp)://hostname:port diff --git a/tests/infra/network.py b/tests/infra/network.py index 247527646e80..9808235df9e1 100644 --- a/tests/infra/network.py +++ b/tests/infra/network.py @@ -662,24 +662,18 @@ def start(self, args, **kwargs): self.consortium.update_recovery_threshold_from_node(primary) def open(self, args): - def get_target_node(args, primary): - # HTTP/2 does not currently support forwarding - if args.http2: - return primary - return self.find_random_node() - primary, _ = self.find_primary() - self.consortium.activate(get_target_node(args, primary)) + self.consortium.activate(self.find_random_node()) if args.js_app_bundle: self.consortium.set_js_app_from_dir( - remote_node=get_target_node(args, primary), + remote_node=self.find_random_node(), bundle_path=args.js_app_bundle, ) for path in args.jwt_issuer: self.consortium.set_jwt_issuer( - remote_node=get_target_node(args, primary), json_path=path + remote_node=self.find_random_node(), json_path=path ) if self.jwt_issuer: @@ -691,7 +685,7 @@ def get_target_node(args, primary): self.create_users(initial_users, args.participants_curve) self.consortium.add_users_and_transition_service_to_open( - get_target_node(args, primary), initial_users + self.find_random_node(), initial_users ) self.status = ServiceStatus.OPEN LOG.info(f"Initial set of users added: {len(initial_users)}") diff --git a/tests/infra/node.py b/tests/infra/node.py index b7894c24fc83..eaf63d5fd928 100644 --- a/tests/infra/node.py +++ b/tests/infra/node.py @@ -1063,9 +1063,10 @@ def _client( akwargs["protocol"] = ( kwargs.get("protocol") if "protocol" in kwargs else "https" ) - if rpc_interface.app_protocol == "HTTP2": - akwargs["http1"] = False - akwargs["http2"] = True + if rpc_interface.app_protocol != "HTTP1": + raise ValueError( + f"Python clients only support HTTP/1.1 interfaces, but {interface_name} is {rpc_interface.app_protocol}" + ) akwargs.update(self.session_auth(identity)) akwargs.update(self.signing_auth(signing_identity)) diff --git a/tests/limits.py b/tests/limits.py index 536f9046ceab..fc9a6850c690 100644 --- a/tests/limits.py +++ b/tests/limits.py @@ -132,14 +132,12 @@ def run_transaction_size_limit_checks(args): if __name__ == "__main__": cr = ConcurrentRunner() - if not cr.args.http2: - # No support for forwarding with HTTP/2 - cr.add( - "parser_limits", - run_parser_limits_checks, - package="samples/apps/logging/logging", - nodes=infra.e2e_args.max_nodes(cr.args, f=0), - ) + cr.add( + "parser_limits", + run_parser_limits_checks, + package="samples/apps/logging/logging", + nodes=infra.e2e_args.max_nodes(cr.args, f=0), + ) cr.add( "transaction_size_limit", diff --git a/tests/memberclient.py b/tests/memberclient.py index d1d39e3e7060..b7f429111cd8 100644 --- a/tests/memberclient.py +++ b/tests/memberclient.py @@ -41,8 +41,8 @@ def test_missing_signature_header(network, args): def make_signature_corrupter(fn): class SignatureCorrupter(infra.clients.HttpSig): - def auth_flow(self, request): - yield fn(next(super().auth_flow(request))) + def __call__(self, request): + return fn(super().__call__(request)) return SignatureCorrupter diff --git a/tests/partitions_test.py b/tests/partitions_test.py index 8e789e997393..dfd6689aad9e 100644 --- a/tests/partitions_test.py +++ b/tests/partitions_test.py @@ -824,7 +824,6 @@ def blocking_send(): "Session consistency is provided, and inconsistencies after elections are replaced by errors" ) @reqs.supports_methods("/app/log/public") -@reqs.no_http2() def test_session_consistency(network, args): # Ensure we have 5 nodes original_size = network.resize(5, args) @@ -1595,9 +1594,7 @@ def run_forwarding_and_sessions(args): with partitioned_network(args) as network: test_forwarding_timeout(network, args) test_invalidated_blocking_calls(network, args) - # HTTP2 doesn't support forwarding - if not args.http2: - test_session_consistency(network, args) + test_session_consistency(network, args) def run_recovery_elections(args): diff --git a/tests/requirements.txt b/tests/requirements.txt index 490d5615586d..1092b745b9c3 100644 --- a/tests/requirements.txt +++ b/tests/requirements.txt @@ -8,13 +8,7 @@ GitPython >= 3.1.45, < 4 better_exceptions >= 0.3.3, < 0.4 pyasn1 >= 0.6.1, < 0.7 Jinja2 >= 3.1.6, < 4 -# Later versions of httpx cause failures in long forwarding, -# extended character range and JWT tests. -httpx[http2] == 0.23.* -# infra.clients.HttpxClient depends on private httpcore internals (module -# path, class, method and pool attribute names) to set TCP_NODELAY; pin -# exactly, since these are not covered by httpcore's public API guarantees. -httpcore == 0.16.3 +requests >= 2.32.5, < 3 locust >= 2.41.6, < 3 JWCrypto >= 1.5.6, < 2 rich >= 14.2.0, < 15 diff --git a/tests/suite/test_requirements.py b/tests/suite/test_requirements.py index e00955bb6dec..2eac2b7fedcc 100644 --- a/tests/suite/test_requirements.py +++ b/tests/suite/test_requirements.py @@ -153,15 +153,6 @@ def check(network, args, *nargs, **kwargs): return ensure_reqs(check) -def no_http2(): - # HTTP/2 does not support forwarding - def check(network, args, *nargs, **kwargs): - if args.http2: - raise TestRequirementsNotMet("Test not run with HTTP/2") - - return ensure_reqs(check) - - def snp_only(): def check(*args, **kwargs): if not SNP_SUPPORT: diff --git a/tests/tvc.py b/tests/tvc.py index 49ff76c07a4d..df98460cd1a9 100644 --- a/tests/tvc.py +++ b/tests/tvc.py @@ -5,7 +5,7 @@ import json import random -import httpx +import requests """ 1. Run sandbox @@ -31,6 +31,7 @@ KEY = "0" VALUE = "value" +TIMEOUT_S = 5 def log(**kwargs): @@ -52,14 +53,15 @@ def retry(call, urls, **kwargs): while response is None or response.status_code not in (200, 204): try: url = random.choice(urls) - response = call(url, **kwargs) - except (httpx.ReadTimeout, httpx.ConnectTimeout): + response = call(url, timeout=TIMEOUT_S, **kwargs) + except requests.exceptions.Timeout: pass return response def run(targets, cacert): - session = httpx.Client(verify=cacert) + session = requests.Session() + session.verify = cacert tx = -1 key_urls = [f"{target}/records/{KEY}" for target in targets] while True: