307 lines
10 KiB
Python
307 lines
10 KiB
Python
"""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"{}")
|
|
rejected._control = lambda: (_ for _ in ()).throw(broker.PolicyError("no"))
|
|
rejected.do_POST()
|
|
assert json.loads(rejected.wfile.getvalue()) == {"error": "request rejected"}
|
|
|
|
|
|
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"
|