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: