atlas-iac/testing/tests/test_hermes_scm_broker_handler_coverage.py

310 lines
11 KiB
Python
Raw Permalink Normal View History

"""Behavioral branch coverage for SCM broker request handling."""
from __future__ import annotations
import io
import json
import sys
from email.message import Message
import pytest
from testing.tests.test_hermes_scm_broker_support import _load, _receive_command
class _Connection:
def __init__(self):
self.timeouts = []
def settimeout(self, value):
self.timeouts.append(value)
def _headers(*, content_type="application/json", body=b""):
value = Message()
value["Content-Type"] = content_type
value["Content-Length"] = str(len(body))
return value
def _handler(broker, *, path="/healthz", body=b"", content_type="application/json"):
handler = object.__new__(broker.BrokerHandler)
handler.path = path
handler.headers = _headers(content_type=content_type, body=body)
handler.rfile = io.BytesIO(body)
handler.wfile = io.BytesIO()
handler.connection = _Connection()
handler.events = []
handler.send_response = lambda status: handler.events.append(("status", status))
handler.send_header = lambda name, value: handler.events.append((name, value))
handler.end_headers = lambda: handler.events.append(("headers", "done"))
return handler
def test_header_validation_accepts_safe_ascii_and_rejects_controls():
broker = _load("scm_broker")
handler = _handler(broker)
handler._validate_headers()
class HeaderBag:
def items(self):
return [("X-Test", "bad\x01value")]
def get_all(self, _name, default):
return default
handler.headers = HeaderBag()
with pytest.raises(broker.PolicyError, match="controls"):
handler._validate_headers()
class UnicodeHeader(HeaderBag):
def items(self):
return [("X-Test", "máin")]
handler.headers = UnicodeHeader()
with pytest.raises(broker.PolicyError, match="ASCII"):
handler._validate_headers()
def test_json_reject_and_stream_response_methods(monkeypatch):
broker = _load("scm_broker")
handler = _handler(broker)
handler._json(202, b'{"ok":true}')
assert ("status", 202) in handler.events
assert handler.wfile.getvalue() == b'{"ok":true}'
handler = _handler(broker)
handler._reject(403)
assert ("status", 403) in handler.events
assert json.loads(handler.wfile.getvalue()) == {"error": "request rejected"}
handler = _handler(broker)
handler._stream(200, "application/test", io.BytesIO(b"streamed"), 8)
assert handler.wfile.getvalue() == b"streamed"
assert handler.connection.timeouts
handler.log_message("ignored %s", "value")
handler = _handler(broker)
ticks = iter((0.0, 121.0))
monkeypatch.setattr(broker.time, "monotonic", lambda: next(ticks))
with pytest.raises(broker.PolicyError, match="write deadline"):
handler._stream(200, "application/test", io.BytesIO(b"x"), 1)
def test_get_health_and_git_discovery_paths(monkeypatch):
broker = _load("scm_broker")
handler = _handler(broker)
handler.do_GET()
assert json.loads(handler.wfile.getvalue()) == {"status": "ok"}
closed = []
class Body(io.BytesIO):
def close(self):
closed.append(True)
super().close()
handler = _handler(
broker,
path="/git/atlas/cassandra.git/info/refs?service=git-upload-pack",
)
monkeypatch.setattr(broker, "read_token", lambda: "sentinel")
monkeypatch.setattr(
broker,
"_upstream_git_request",
lambda *args, **kwargs: (Body(b"advertisement"), 13),
)
handler.do_GET()
assert handler.wfile.getvalue() == b"advertisement"
assert closed == [True]
handler = _handler(broker, path="/git/atlas/cassandra.git/git-upload-pack")
handler.do_GET()
assert json.loads(handler.wfile.getvalue()) == {"error": "request rejected"}
def test_post_dispatches_control_and_git_or_rejects(monkeypatch):
broker = _load("scm_broker")
control = _handler(broker, path="/v1/metadata", body=b"{}")
git = _handler(broker, path="/git/atlas/cassandra.git/git-upload-pack", body=b"")
calls = []
control._control = lambda: calls.append("control")
control._git_rpc = lambda: calls.append("git")
git._control = lambda: calls.append("control")
git._git_rpc = lambda: calls.append("git")
control.do_POST()
git.do_POST()
assert calls == ["control", "git"]
rejected = _handler(broker, path="/v1/metadata", body=b"{}")
events = []
monkeypatch.setattr(broker, "log_rejection", lambda phase, error: events.append((phase, type(error).__name__)))
rejected._control = lambda: (_ for _ in ()).throw(broker.PolicyError("no"))
rejected.do_POST()
assert json.loads(rejected.wfile.getvalue()) == {"error": "request rejected"}
assert events == [("control", "PolicyError")]
def test_control_metadata_and_draft_fields(monkeypatch):
broker = _load("scm_broker")
monkeypatch.setattr(broker, "read_token", lambda: "sentinel")
monkeypatch.setattr(
broker, "read", lambda path, token: json.dumps({"path": path}).encode()
)
monkeypatch.setattr(
broker,
"create_draft",
lambda token, **data: json.dumps({"token_used": bool(token), **data}).encode(),
)
metadata_body = b'{"path":"/api/v1/repos/titan/cassandra"}'
metadata = _handler(broker, path="/v1/metadata", body=metadata_body)
metadata._control()
assert json.loads(metadata.wfile.getvalue())["path"].endswith("cassandra")
draft_data = {
"base": "main",
"body": "Review evidence",
"head": "feature/coverage",
"head_sha": "a" * 40,
"repo": "cassandra",
"title": "WIP: Coverage",
}
draft_body = json.dumps(draft_data).encode()
draft = _handler(broker, path="/v1/drafts", body=draft_body)
draft._control()
assert json.loads(draft.wfile.getvalue())["repo"] == "cassandra"
for path, body, match in (
("/v1/metadata", b'{"wrong":"field"}', "fields"),
("/v1/drafts", b'{"repo":1}', "fields"),
):
with pytest.raises(broker.PolicyError, match=match):
_handler(broker, path=path, body=body)._control()
monkeypatch.setattr(broker, "read", lambda *_a, **_k: b"sentinel")
with pytest.raises(broker.PolicyError, match="reflected"):
_handler(broker, path="/v1/metadata", body=metadata_body)._control()
@pytest.mark.parametrize("service", ["git-upload-pack", "git-receive-pack"])
def test_git_rpc_streams_allowed_upload_and_receive(monkeypatch, service):
broker = _load("scm_broker")
if service == "git-receive-pack":
body = _receive_command(b"0" * 40, b"1" * 40, b"refs/heads/hermes/coverage")
else:
body = b"upload-request"
handler = _handler(
broker,
path=f"/git/atlas/cassandra.git/{service}",
body=body,
content_type=f"application/x-{service}-request",
)
monkeypatch.setattr(broker, "read_token", lambda: "sentinel")
seen = []
def upstream(target, **kwargs):
seen.append((target, kwargs))
return io.BytesIO(b"result"), 6
monkeypatch.setattr(broker, "_upstream_git_request", upstream)
handler._git_rpc()
assert handler.wfile.getvalue() == b"result"
assert seen[0][0].endswith(service)
assert seen[0][1]["body_length"] == len(body)
def test_signed_receive_pack_advances_ledger_only_after_exact_live_head(monkeypatch):
broker = _load("scm_broker")
body = _receive_command(b"0" * 40, b"1" * 40, b"refs/heads/hermes/coverage")
handler = _handler(
broker, path="/git/atlas/cassandra.git/git-receive-pack", body=body,
content_type="application/x-git-receive-pack-request",
)
handler.headers["X-Hermes-Task-Grant"] = "bounded-signed-grant"
claims = {"repo": "cassandra", "ref": "hermes/coverage", "expected_old": "0" * 40, "new_head": "1" * 40}
calls = []
class Ledger:
def authorize_update(self, value):
calls.append(("authorize", value))
def commit(self, value):
calls.append(("commit", value))
monkeypatch.setattr(broker, "read_token", lambda: "sentinel")
monkeypatch.setattr(broker, "verify_grant", lambda value: calls.append(("header", value)) or claims)
monkeypatch.setattr(broker, "_task_ledger", lambda: Ledger())
monkeypatch.setattr(broker, "_branch_head", lambda *_args: "1" * 40)
monkeypatch.setattr(broker, "_upstream_git_request", lambda *_args, **_kwargs: (io.BytesIO(b"result"), 6))
handler._git_rpc()
assert calls == [("header", "bounded-signed-grant"), ("authorize", claims), ("commit", claims)]
calls.clear()
handler = _handler(
broker, path="/git/atlas/cassandra.git/git-receive-pack", body=body,
content_type="application/x-git-receive-pack-request",
)
handler.headers["X-Hermes-Task-Grant"] = "bounded-signed-grant"
monkeypatch.setattr(broker, "_branch_head", lambda *_args: "0" * 40)
with pytest.raises(broker.PolicyError, match="did not advance"):
handler._git_rpc()
assert not [call for call in calls if call[0] == "commit"]
@pytest.mark.parametrize(
("path", "content_type", "transfer", "match"),
[
(
"/git/atlas/cassandra.git/info/refs?service=git-upload-pack",
"application/x-git-upload-pack-request",
None,
"operation",
),
(
"/git/atlas/cassandra.git/git-upload-pack",
"text/plain",
None,
"request type",
),
(
"/git/atlas/cassandra.git/git-upload-pack",
"application/x-git-upload-pack-request",
"chunked",
"request type",
),
],
)
def test_git_rpc_rejects_discovery_wrong_type_and_chunking(
path, content_type, transfer, match
):
broker = _load("scm_broker")
handler = _handler(broker, path=path, content_type=content_type)
if transfer:
handler.headers["Transfer-Encoding"] = transfer
with pytest.raises(broker.PolicyError, match=match):
handler._git_rpc()
def test_broker_main_constructs_bounded_server(monkeypatch):
broker = _load("scm_broker")
seen = []
class Server:
def __init__(self, address, handler):
seen.append((address, handler))
def serve_forever(self):
seen.append("served")
monkeypatch.setattr(broker, "BoundedThreadingHTTPServer", Server)
monkeypatch.setattr(broker, "seed_adoptions", lambda *_args: None)
monkeypatch.setattr(broker, "_task_ledger", lambda: object())
monkeypatch.setattr(broker, "read_token", lambda: "token")
monkeypatch.setattr(
sys, "argv", ["broker", "--listen", "127.0.0.1", "--port", "9191"]
)
assert broker.main() == 0
assert seen[0] == (("127.0.0.1", 9191), broker.BrokerHandler)
assert seen[1] == "served"