Files
agent-call/tests/test_sip_management.py

457 lines
17 KiB
Python

from __future__ import annotations
import json
import threading
import unittest
from http.client import HTTPConnection
from pathlib import Path
from tempfile import TemporaryDirectory
from typing import Any
from urllib.parse import urlsplit
from agent_call.sip_management import (
ADMIN_AUDIENCE,
READ_AUDIENCE,
ConfigurationError,
Principal,
SipManagementService,
load_token_map,
make_server,
validate_token_separation,
)
ROOT = Path(__file__).resolve().parents[1]
ADMIN_BEARER = "admin-test"
READ_BEARER = "saas-test"
class SipManagementTests(unittest.TestCase):
def test_codec_mapping_and_token_domains_are_separate(self) -> None:
admin = load_token_map(
json.dumps(
{
ADMIN_BEARER: {
"subject": "ops",
"issuer": "ops-issuer",
"audience": ADMIN_AUDIENCE,
"scopes": ["*"],
"trunk_ids": "*",
}
}
),
"admin",
)
read = load_token_map(
json.dumps(
{
READ_BEARER: {
"subject": "saas",
"issuer": "saas-issuer",
"audience": READ_AUDIENCE,
"scopes": ["sip.trunk.read"],
"trunk_ids": ["trunk-a"],
}
}
),
"read",
)
validate_token_separation(admin, read, {"scheduler-test": {}})
self.assertNotEqual(set(admin), set(read))
with self.assertRaises(ConfigurationError):
validate_token_separation(admin, read, {ADMIN_BEARER: {}})
def test_persistent_versioned_publish_readonly_and_rollback(self) -> None:
with TemporaryDirectory() as directory:
service = SipManagementService(Path(directory) / "sip.sqlite3")
try:
service.upsert_cell(
"cell-a",
{
"egress_pool_id": "egress-a",
"codec_capabilities": ["PCMA", "PCMU"],
"status": "healthy",
"max_concurrency": 100,
},
expected_revision=0,
actor="ops",
request_id="cell-1",
)
created = service.upsert_trunk(
"trunk-a",
self._trunk_payload(["PCMA"]),
expected_revision=0,
actor="ops",
request_id="trunk-1",
)
self.assertEqual(created["latest_revision"], 1)
published = service.publish_trunk(
"trunk-a", expected_revision=1, actor="ops", request_id="publish-1"
)
self.assertEqual(published["status"], "published")
self.assertEqual(published["compatible_cell_ids"], ["cell-a"])
self.assertEqual(
service.list_publications("trunk-a")[0]["status"], "pending"
)
readonly = service.get_readonly_trunk("trunk-a", frozenset({"trunk-a"}))
self.assertEqual(
readonly["config"]["codec_profile"]["allowed"], ["PCMA"]
)
self.assertNotIn("credential_ref", readonly["config"]["sip"])
updated = service.upsert_trunk(
"trunk-a",
self._trunk_payload(["PCMU"]),
expected_revision=1,
actor="ops",
request_id="trunk-2",
)
self.assertEqual(updated["latest_revision"], 2)
with self.assertRaisesRegex(Exception, "revision changed"):
service.publish_trunk(
"trunk-a",
expected_revision=1,
actor="ops",
request_id="publish-old",
)
service.publish_trunk(
"trunk-a", expected_revision=2, actor="ops", request_id="publish-2"
)
rolled_back = service.rollback_trunk(
"trunk-a",
1,
expected_revision=2,
actor="ops",
request_id="rollback-1",
)
self.assertEqual(rolled_back["active_revision"], 3)
self.assertEqual(
rolled_back["active"]["codec_profile"]["allowed"], ["PCMA"]
)
self.assertEqual(
[entry["action"] for entry in service.list_audit("trunk-a")],
["upsert", "publish", "upsert", "publish", "rollback"],
)
finally:
service.close()
def test_state_survives_restart_and_secret_refs_are_not_exposed(self) -> None:
with TemporaryDirectory() as directory:
database = Path(directory) / "sip.sqlite3"
service = SipManagementService(database)
service.upsert_cell(
"cell-a",
{
"egress_pool_id": "egress-a",
"codec_capabilities": ["PCMA"],
"status": "healthy",
"max_concurrency": 10,
},
expected_revision=0,
actor="ops",
request_id="cell-1",
)
service.upsert_trunk(
"trunk-a",
{
**self._trunk_payload(["PCMA"]),
"sip": {
**self._trunk_payload(["PCMA"])["sip"],
"auth_mode": "digest",
"credential_ref": "secret://provider-a",
},
},
expected_revision=0,
actor="ops",
request_id="trunk-1",
)
service.publish_trunk(
"trunk-a", expected_revision=1, actor="ops", request_id="publish-1"
)
service.close()
reopened = SipManagementService(database)
try:
admin = reopened.get_trunk("trunk-a")
readonly = reopened.get_readonly_trunk(
"trunk-a", frozenset({"trunk-a"})
)
self.assertTrue(admin["active"]["credential_configured"])
self.assertNotIn("credential_ref", admin["active"]["sip"])
self.assertNotIn("credential_ref", readonly["config"]["sip"])
finally:
reopened.close()
def test_openapi_separates_admin_and_saas_readonly_surfaces(self) -> None:
contract = (ROOT / "docs/contracts/sip-management.openapi.yaml").read_text(
encoding="utf-8"
)
self.assertIn("SipAdminBearer", contract)
self.assertIn("SaasTrunkReadBearer", contract)
self.assertIn("/admin/v1/trunks/{trunk_id}/publish:", contract)
self.assertIn("/readonly/v1/sip/trunks:", contract)
self.assertNotIn("HTTP_TOKENS", contract)
def test_publish_rejects_without_compatible_cell(self) -> None:
with TemporaryDirectory() as directory:
service = SipManagementService(Path(directory) / "sip.sqlite3")
try:
service.upsert_cell(
"cell-pcmu",
{
"egress_pool_id": "egress-a",
"codec_capabilities": ["PCMU"],
"status": "healthy",
"max_concurrency": 10,
},
expected_revision=0,
actor="ops",
request_id="cell-1",
)
service.upsert_trunk(
"trunk-a",
self._trunk_payload(["PCMA"]),
expected_revision=0,
actor="ops",
request_id="trunk-1",
)
with self.assertRaisesRegex(Exception, "no healthy Cell"):
service.publish_trunk(
"trunk-a",
expected_revision=1,
actor="ops",
request_id="publish-1",
)
finally:
service.close()
def test_http_readonly_token_cannot_write_or_read_admin_api(self) -> None:
with TemporaryDirectory() as directory:
service = SipManagementService(Path(directory) / "sip.sqlite3")
admin_tokens = {
ADMIN_BEARER: Principal(
"ops",
"ops-issuer",
ADMIN_AUDIENCE,
frozenset({"*"}),
frozenset({"*"}),
"admin",
)
}
read_tokens = {
READ_BEARER: Principal(
"saas",
"saas-issuer",
READ_AUDIENCE,
frozenset({"sip.trunk.read"}),
frozenset({"trunk-a"}),
"read",
)
}
server = make_server(
service,
"127.0.0.1",
0,
admin_tokens=admin_tokens,
read_tokens=read_tokens,
)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
base = f"http://127.0.0.1:{server.server_port}"
try:
status, body = self._request(
base,
"GET",
"/healthz/live",
bearer=None,
)
self.assertEqual(status, 200)
self.assertEqual(body["mode"], "mock")
status, _ = self._request(
base,
"PUT",
"/admin/v1/trunks/trunk-a",
bearer=READ_BEARER,
headers={"If-Match": "0", "X-Request-ID": "write-1"},
body=self._trunk_payload(["PCMA"]),
)
self.assertEqual(status, 401)
status, _ = self._request(
base,
"GET",
"/readonly/v1/sip/trunks",
bearer=ADMIN_BEARER,
)
self.assertEqual(status, 401)
status, _ = self._request(
base,
"POST",
"/readonly/v1/sip/trunks/trunk-a",
bearer=READ_BEARER,
body={},
)
self.assertEqual(status, 405)
finally:
server.shutdown()
server.server_close()
thread.join(timeout=2)
service.close()
def test_http_admin_publish_and_readonly_uses_active_revision(self) -> None:
with TemporaryDirectory() as directory:
service = SipManagementService(Path(directory) / "sip.sqlite3")
admin_tokens = {
ADMIN_BEARER: Principal(
"ops",
"ops-issuer",
ADMIN_AUDIENCE,
frozenset({"*"}),
frozenset({"*"}),
"admin",
)
}
read_tokens = {
READ_BEARER: Principal(
"saas",
"saas-issuer",
READ_AUDIENCE,
frozenset({"sip.trunk.read"}),
frozenset({"trunk-a"}),
"read",
)
}
server = make_server(
service,
"127.0.0.1",
0,
admin_tokens=admin_tokens,
read_tokens=read_tokens,
)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
base = f"http://127.0.0.1:{server.server_port}"
try:
status, _ = self._request(
base,
"PUT",
"/admin/v1/cells/cell-a",
bearer=ADMIN_BEARER,
headers={"If-Match": "0", "X-Request-ID": "cell-1"},
body={
"egress_pool_id": "egress-a",
"codec_capabilities": ["PCMA", "PCMU"],
"status": "healthy",
"max_concurrency": 10,
},
)
self.assertEqual(status, 201)
status, _ = self._request(
base,
"PUT",
"/admin/v1/trunks/trunk-a",
bearer=ADMIN_BEARER,
headers={"If-Match": "0", "X-Request-ID": "trunk-1"},
body=self._trunk_payload(["PCMA"]),
)
self.assertEqual(status, 201)
status, _ = self._request(
base,
"POST",
"/admin/v1/trunks/trunk-a/publish",
bearer=ADMIN_BEARER,
headers={"If-Match": "1", "X-Request-ID": "publish-1"},
)
self.assertEqual(status, 200)
status, body = self._request(
base,
"GET",
"/readonly/v1/sip/trunks/trunk-a",
bearer=READ_BEARER,
)
self.assertEqual(status, 200)
self.assertEqual(body["revision"], 1)
self.assertEqual(body["config"]["codec_profile"]["allowed"], ["PCMA"])
self.assertNotIn("latest", body)
self.assertNotIn("compatible_cell_ids", body)
status, _ = self._request(
base,
"PUT",
"/admin/v1/trunks/trunk-a",
bearer=ADMIN_BEARER,
headers={"If-Match": "1", "X-Request-ID": "trunk-2"},
body=self._trunk_payload(["PCMU"]),
)
self.assertEqual(status, 200)
status, body = self._request(
base,
"GET",
"/readonly/v1/sip/trunks/trunk-a",
bearer=READ_BEARER,
)
self.assertEqual(status, 200)
self.assertEqual(body["revision"], 1)
self.assertEqual(body["config"]["codec_profile"]["allowed"], ["PCMA"])
finally:
server.shutdown()
server.server_close()
thread.join(timeout=2)
service.close()
@staticmethod
def _trunk_payload(codecs: list[str]) -> dict[str, Any]:
return {
"display_name": "Provider A",
"enabled": True,
"sip": {
"host": "61.132.228.221",
"port": 5060,
"transport": "udp",
"auth_mode": "ip",
"register": False,
},
"codec_profile": {"allowed": codecs, "preferred": codecs[0]},
"caller_ids": ["BD93205882"],
"dial_prefix": "7089",
"egress_pool_id": "egress-a",
"max_concurrency": 100,
"max_cps": 10,
}
@staticmethod
def _request(
base: str,
method: str,
path: str,
*,
bearer: str | None,
headers: dict[str, str] | None = None,
body: Any | None = None,
) -> tuple[int, dict[str, Any]]:
request_headers = {"Accept": "application/json"}
if bearer:
request_headers["Authorization"] = f"Bearer {bearer}"
if headers:
request_headers.update(headers)
data = None
if body is not None:
data = json.dumps(body).encode("utf-8")
request_headers["Content-Type"] = "application/json"
parsed = urlsplit(base)
host = parsed.hostname
port = parsed.port
if parsed.scheme != "http" or host != "127.0.0.1" or port is None:
raise AssertionError("test helper only permits a local HTTP server")
connection = HTTPConnection(host, port, timeout=3)
try:
connection.request(method, path, body=data, headers=request_headers)
response = connection.getresponse()
return response.status, json.loads(response.read())
finally:
connection.close()
if __name__ == "__main__":
unittest.main()