Hotspots connect on PROXY.LISTEN_PORT with in-process inject into SYSTEM, session-preserving SIGHUP reload, and composition-root wiring in main.pull/1/head
parent
ca59b886fb
commit
614c882bcb
@ -1,10 +1,14 @@
|
||||
"""Proxy application layer (Phase 3)."""
|
||||
|
||||
from .deployment import is_proxy_inject_only, normalize_proxy_target, proxy_target_system
|
||||
from .packet_helpers import peer_id_from_packet
|
||||
from .use_cases import ProxySlotError, ProxyUseCases
|
||||
|
||||
__all__ = [
|
||||
"ProxySlotError",
|
||||
"ProxyUseCases",
|
||||
"is_proxy_inject_only",
|
||||
"normalize_proxy_target",
|
||||
"peer_id_from_packet",
|
||||
"proxy_target_system",
|
||||
]
|
||||
|
||||
@ -0,0 +1,43 @@
|
||||
"""Proxy deployment policy from config dict (no I/O; used at startup/reload)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
|
||||
def proxy_target_system(config: dict[str, Any]) -> str | None:
|
||||
proxy = config.get("PROXY", {})
|
||||
target = proxy.get("TARGET_SYSTEM")
|
||||
return str(target) if target else None
|
||||
|
||||
|
||||
def config_has_enabled_master(config: dict[str, Any]) -> bool:
|
||||
"""True when config defines at least one enabled MASTER (adn-server, not parrot-only)."""
|
||||
systems = config.get("SYSTEMS", {})
|
||||
if not isinstance(systems, dict):
|
||||
return False
|
||||
return any(
|
||||
isinstance(cfg, dict) and cfg.get("ENABLED", True) and cfg.get("MODE") == "MASTER"
|
||||
for cfg in systems.values()
|
||||
)
|
||||
|
||||
|
||||
def is_proxy_inject_only(config: dict[str, Any], system_name: str) -> bool:
|
||||
target = proxy_target_system(config)
|
||||
return target is not None and target == system_name
|
||||
|
||||
|
||||
def normalize_proxy_target(config: dict[str, Any]) -> None:
|
||||
"""Strip direct UDP bind fields from inject-only proxy target (D-23)."""
|
||||
target = proxy_target_system(config)
|
||||
if not target:
|
||||
return
|
||||
sys_cfg = config.get("SYSTEMS", {}).get(target)
|
||||
if not isinstance(sys_cfg, dict):
|
||||
return
|
||||
port = sys_cfg.pop("PORT", None)
|
||||
sys_cfg.pop("IP", None)
|
||||
if port is not None:
|
||||
sys_cfg["_REPORT_BASE_PORT"] = int(port)
|
||||
else:
|
||||
sys_cfg.setdefault("_REPORT_BASE_PORT", 56400)
|
||||
@ -0,0 +1,17 @@
|
||||
"""Wire packets for proxy session teardown (legacy reaper parity)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
# Homebrew command prefixes (wire vocabulary; no I/O)
|
||||
_MSTCL = b"MSTCL"
|
||||
_RPTCL = b"RPTCL"
|
||||
|
||||
CLIENT_TEARDOWN_REPEAT = 3
|
||||
|
||||
|
||||
def master_teardown_packet(peer_id: bytes) -> bytes:
|
||||
return _RPTCL + peer_id
|
||||
|
||||
|
||||
def client_teardown_packet() -> bytes:
|
||||
return _MSTCL
|
||||
@ -1,9 +1,29 @@
|
||||
"""Proxy infrastructure adapters (Phase 3)."""
|
||||
|
||||
from .config import apply_proxy_env_overrides, proxy_settings
|
||||
from .hbp_adapters import FanInClientSender, HbpMasterPeerRegistry, InProcessHbpSink
|
||||
from .ip_blacklist import InMemoryProxyIpBlacklist
|
||||
from .reply_transport import ProxyReplyTransport
|
||||
from .runtime import ProxyServiceState, apply_proxy_config_reload, start_proxy_service
|
||||
from .rpto_queue import InMemoryPendingRptoQueue
|
||||
from .session_executor import apply_session_teardown
|
||||
from .slot_store import InMemoryProxySlotStore
|
||||
from .udp_fanin import ProxyFanInProtocol, listen_proxy_fanin
|
||||
|
||||
__all__ = [
|
||||
"FanInClientSender",
|
||||
"HbpMasterPeerRegistry",
|
||||
"InMemoryPendingRptoQueue",
|
||||
"InMemoryProxyIpBlacklist",
|
||||
"InMemoryProxySlotStore",
|
||||
"InProcessHbpSink",
|
||||
"ProxyFanInProtocol",
|
||||
"ProxyReplyTransport",
|
||||
"ProxyServiceState",
|
||||
"apply_proxy_config_reload",
|
||||
"apply_proxy_env_overrides",
|
||||
"apply_session_teardown",
|
||||
"listen_proxy_fanin",
|
||||
"proxy_settings",
|
||||
"start_proxy_service",
|
||||
]
|
||||
|
||||
@ -0,0 +1,41 @@
|
||||
"""PROXY runtime settings from config dict (infrastructure; no business rules)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
|
||||
def apply_proxy_env_overrides(config: dict[str, Any]) -> None:
|
||||
"""Apply ADN_PROXY_* environment overrides (design §5.6)."""
|
||||
import os
|
||||
|
||||
proxy = config.setdefault("PROXY", {})
|
||||
if os.environ.get("ADN_PROXY_DEBUG", "").strip() in ("1", "true", "TRUE", "yes", "YES"):
|
||||
proxy["DEBUG"] = True
|
||||
listen_port = os.environ.get("ADN_PROXY_LISTENPORT", "").strip()
|
||||
if listen_port:
|
||||
proxy["LISTEN_PORT"] = int(listen_port)
|
||||
if os.environ.get("ADN_PROXY_IPV6", "").strip() in ("1", "true", "TRUE", "yes", "YES"):
|
||||
proxy["LISTEN_IP"] = "::"
|
||||
|
||||
|
||||
def proxy_settings(config: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Resolved PROXY runtime settings with defaults."""
|
||||
proxy = config.get("PROXY", {})
|
||||
black_list = proxy.get("BLACK_LIST") or []
|
||||
if not isinstance(black_list, list):
|
||||
black_list = []
|
||||
ip_black_list = proxy.get("IP_BLACK_LIST") or {}
|
||||
if not isinstance(ip_black_list, dict):
|
||||
ip_black_list = {}
|
||||
return {
|
||||
"listen_port": int(proxy.get("LISTEN_PORT", 62031)),
|
||||
"listen_ip": str(proxy.get("LISTEN_IP") or ""),
|
||||
"target_system": str(proxy.get("TARGET_SYSTEM") or ""),
|
||||
"timeout": float(proxy.get("TIMEOUT", 30)),
|
||||
"debug": bool(proxy.get("DEBUG")),
|
||||
"client_info": bool(proxy.get("CLIENT_INFO", True)),
|
||||
"black_list": tuple(int(x) for x in black_list),
|
||||
"ip_black_list": {str(k): float(v) for k, v in ip_black_list.items()},
|
||||
"stats": bool(proxy.get("STATS")),
|
||||
}
|
||||
@ -0,0 +1,47 @@
|
||||
"""HBP adapters for proxy application ports."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any, Protocol
|
||||
|
||||
from adn_server.application.ports import MasterPeerRegistry, ProxyClientSender, ProxyMasterSink
|
||||
from adn_server.domain.proxy import ClientEndpoint
|
||||
|
||||
|
||||
class _MasterHbpReceiver(Protocol):
|
||||
def _master_datagram_received(self, data: bytes, sockaddr: tuple[str, int]) -> None:
|
||||
...
|
||||
|
||||
|
||||
class InProcessHbpSink(ProxyMasterSink):
|
||||
"""Deliver client datagrams to the target MASTER without a UDP hop."""
|
||||
|
||||
def __init__(self, hbp: _MasterHbpReceiver) -> None:
|
||||
self._hbp = hbp
|
||||
|
||||
def inject(self, data: bytes, client_addr: tuple[str, int]) -> None:
|
||||
self._hbp._master_datagram_received(data, client_addr)
|
||||
|
||||
|
||||
class FanInClientSender(ProxyClientSender):
|
||||
"""Send to hotspots through the fan-in UDP transport."""
|
||||
|
||||
def __init__(self, transport: Any) -> None:
|
||||
self._transport = transport
|
||||
|
||||
def send_to_client(self, data: bytes, client: ClientEndpoint) -> None:
|
||||
if self._transport is None:
|
||||
return
|
||||
self._transport.write(data, (client.host, client.port))
|
||||
|
||||
|
||||
class HbpMasterPeerRegistry(MasterPeerRegistry):
|
||||
"""Remove timed-out peers from MASTER ``_peers``."""
|
||||
|
||||
def __init__(self, hbp: Any) -> None:
|
||||
self._hbp = hbp
|
||||
|
||||
def remove_peer(self, peer_id: bytes) -> None:
|
||||
peers = getattr(self._hbp, "_peers", None)
|
||||
if isinstance(peers, dict):
|
||||
peers.pop(peer_id, None)
|
||||
@ -0,0 +1,21 @@
|
||||
"""In-memory IP blacklist for proxy fan-in."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from adn_server.application.ports import ProxyIpBlacklist
|
||||
|
||||
|
||||
class InMemoryProxyIpBlacklist(ProxyIpBlacklist):
|
||||
def __init__(self, initial: dict[str, float] | None = None) -> None:
|
||||
self._entries: dict[str, float] = dict(initial or {})
|
||||
|
||||
def block_until(self, host: str, expire_at: float) -> None:
|
||||
self._entries[host] = expire_at
|
||||
|
||||
def is_blocked(self, host: str, now: float) -> bool:
|
||||
expire = self._entries.get(host)
|
||||
return expire is not None and now < expire
|
||||
|
||||
def merge_static_entries(self, entries: dict[str, float]) -> None:
|
||||
"""Apply config IP_BLACK_LIST entries (runtime PRBL blocks are kept)."""
|
||||
self._entries.update(entries)
|
||||
@ -0,0 +1,32 @@
|
||||
"""Route MASTER HBP replies through the proxy fan-in UDP socket."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Callable
|
||||
from typing import Protocol
|
||||
|
||||
from adn_server.infrastructure.hbp_constants import PRBL
|
||||
|
||||
|
||||
class _DatagramWriter(Protocol):
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
...
|
||||
|
||||
|
||||
class ProxyReplyTransport:
|
||||
"""Wrap the fan-in transport so MASTER replies leave via LISTEN_PORT."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
fanin_transport: _DatagramWriter,
|
||||
*,
|
||||
prbl_handler: Callable[[bytes, tuple[str, int]], None] | None = None,
|
||||
) -> None:
|
||||
self._fanin = fanin_transport
|
||||
self._prbl_handler = prbl_handler
|
||||
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
if len(data) >= 4 and data[:4] == PRBL and self._prbl_handler is not None:
|
||||
self._prbl_handler(data, addr)
|
||||
return
|
||||
self._fanin.write(data, addr)
|
||||
@ -0,0 +1,225 @@
|
||||
"""Wire integrated hotspot proxy at startup (composition root wiring)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any
|
||||
|
||||
from twisted.internet import reactor
|
||||
from twisted.internet.interfaces import IDelayedCall
|
||||
|
||||
from adn_server.application.proxy import ProxyUseCases
|
||||
from adn_server.application.ports import ProxyClientSender, ProxyMasterSink
|
||||
from adn_server.domain.value_objects import int_id
|
||||
from adn_server.infrastructure.proxy.config import proxy_settings
|
||||
from adn_server.infrastructure.proxy.hbp_adapters import (
|
||||
FanInClientSender,
|
||||
HbpMasterPeerRegistry,
|
||||
InProcessHbpSink,
|
||||
)
|
||||
from adn_server.infrastructure.proxy.ip_blacklist import InMemoryProxyIpBlacklist
|
||||
from adn_server.infrastructure.proxy.rpto_queue import InMemoryPendingRptoQueue
|
||||
from adn_server.infrastructure.proxy.session_executor import apply_session_teardown
|
||||
from adn_server.infrastructure.proxy.slot_store import InMemoryProxySlotStore
|
||||
from adn_server.infrastructure.proxy.udp_fanin import ProxyFanInProtocol, listen_proxy_fanin
|
||||
from adn_server.infrastructure.proxy.reply_transport import ProxyReplyTransport
|
||||
|
||||
|
||||
def _proxy_runtime_snapshot(config: dict[str, Any]) -> dict[str, Any]:
|
||||
settings = proxy_settings(config)
|
||||
target = settings["target_system"]
|
||||
max_peers = int(config.get("SYSTEMS", {}).get(target, {}).get("MAX_PEERS", 1))
|
||||
return {**settings, "max_peers": max_peers}
|
||||
|
||||
|
||||
@dataclass
|
||||
class ProxyServiceState:
|
||||
"""Live proxy handles (for shutdown / reload)."""
|
||||
|
||||
target_system: str
|
||||
use_cases: ProxyUseCases
|
||||
master_sink: ProxyMasterSink
|
||||
client_sender: ProxyClientSender
|
||||
fanin: ProxyFanInProtocol
|
||||
udp_port: Any
|
||||
listen_port: int = 62031
|
||||
listen_ip: str = ""
|
||||
_timers: dict[bytes, IDelayedCall] = field(default_factory=dict)
|
||||
_runtime: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
def stop(self) -> Any:
|
||||
"""Stop timers and UDP listener. Returns Twisted Deferred when a port was bound."""
|
||||
for call in self._timers.values():
|
||||
if call.active():
|
||||
call.cancel()
|
||||
self._timers.clear()
|
||||
if self.udp_port is None:
|
||||
return None
|
||||
port = self.udp_port
|
||||
self.udp_port = None
|
||||
return port.stopListening()
|
||||
|
||||
|
||||
def apply_proxy_config_reload(
|
||||
state: ProxyServiceState,
|
||||
config: dict[str, Any],
|
||||
*,
|
||||
logger: logging.Logger,
|
||||
) -> None:
|
||||
"""Hot-apply PROXY settings on SIGHUP without closing LISTEN_PORT or dropping sessions."""
|
||||
incoming = _proxy_runtime_snapshot(config)
|
||||
bind_changed = (
|
||||
state.listen_port != incoming["listen_port"]
|
||||
or state.listen_ip != incoming["listen_ip"]
|
||||
)
|
||||
target_changed = state.target_system != incoming["target_system"]
|
||||
state._runtime.update(incoming)
|
||||
state.use_cases.apply_runtime_settings(
|
||||
max_peers=incoming["max_peers"],
|
||||
black_list=incoming["black_list"],
|
||||
)
|
||||
ip_bl = state.use_cases._ip_blacklist # noqa: SLF001
|
||||
if isinstance(ip_bl, InMemoryProxyIpBlacklist):
|
||||
ip_bl.merge_static_entries(incoming["ip_black_list"])
|
||||
state.fanin.debug = bool(incoming["debug"])
|
||||
if bind_changed:
|
||||
logger.warning(
|
||||
"(CONFIG-RELOAD) PROXY bind change ignored at runtime "
|
||||
"(still listening on %s:%s); restart adn-server to apply %s:%s",
|
||||
state.listen_ip or "*",
|
||||
state.listen_port,
|
||||
incoming["listen_ip"] or "*",
|
||||
incoming["listen_port"],
|
||||
)
|
||||
if target_changed:
|
||||
logger.warning(
|
||||
"(CONFIG-RELOAD) PROXY TARGET_SYSTEM change ignored at runtime "
|
||||
"(still injecting into %s); restart adn-server to apply %s",
|
||||
state.target_system,
|
||||
incoming["target_system"],
|
||||
)
|
||||
logger.debug(
|
||||
"(CONFIG-RELOAD) proxy settings applied (%s active session(s), port kept open)",
|
||||
len(state.use_cases.list_slots()),
|
||||
)
|
||||
|
||||
|
||||
def start_proxy_service(
|
||||
config: dict[str, Any],
|
||||
protocols: dict[str, Any],
|
||||
*,
|
||||
logger: logging.Logger,
|
||||
) -> ProxyServiceState:
|
||||
"""Start LISTEN_PORT fan-in and inject into ``PROXY.TARGET_SYSTEM`` MASTER."""
|
||||
runtime = _proxy_runtime_snapshot(config)
|
||||
target = runtime["target_system"]
|
||||
target_proto = protocols.get(target)
|
||||
if target_proto is None:
|
||||
raise RuntimeError(f"(PROXY) TARGET_SYSTEM {target!r} has no HBP protocol instance")
|
||||
|
||||
ip_blacklist = InMemoryProxyIpBlacklist(runtime["ip_black_list"])
|
||||
use_cases = ProxyUseCases(
|
||||
InMemoryProxySlotStore(),
|
||||
InMemoryPendingRptoQueue(),
|
||||
max_peers=runtime["max_peers"],
|
||||
black_list=runtime["black_list"],
|
||||
ip_blacklist=ip_blacklist,
|
||||
)
|
||||
master_sink = InProcessHbpSink(target_proto)
|
||||
peer_registry = HbpMasterPeerRegistry(target_proto)
|
||||
fanin = ProxyFanInProtocol(
|
||||
use_cases,
|
||||
master_sink,
|
||||
debug=runtime["debug"],
|
||||
logger=logger,
|
||||
)
|
||||
state = ProxyServiceState(
|
||||
target_system=target,
|
||||
use_cases=use_cases,
|
||||
master_sink=master_sink,
|
||||
client_sender=FanInClientSender(None),
|
||||
fanin=fanin,
|
||||
udp_port=None,
|
||||
listen_port=runtime["listen_port"],
|
||||
listen_ip=runtime["listen_ip"],
|
||||
_runtime=dict(runtime),
|
||||
)
|
||||
|
||||
def _on_client_attached(peer_id: bytes, host: str, port: int, new_session: bool) -> None:
|
||||
rt = state._runtime
|
||||
if (
|
||||
new_session
|
||||
and rt["client_info"]
|
||||
and peer_id != b"\xff\xff\xff\xff"
|
||||
):
|
||||
logger.info(
|
||||
"(PROXY) New client: ID:%s IP:%s Port:%s",
|
||||
str(int_id(peer_id)).rjust(9),
|
||||
host.rjust(15),
|
||||
port,
|
||||
)
|
||||
existing = state._timers.get(peer_id)
|
||||
if existing is not None and existing.active():
|
||||
existing.reset(rt["timeout"])
|
||||
return
|
||||
if existing is not None:
|
||||
existing.cancel()
|
||||
state._timers[peer_id] = reactor.callLater(rt["timeout"], _reap_session, peer_id)
|
||||
|
||||
def _reap_session(peer_id: bytes) -> None:
|
||||
rt = state._runtime
|
||||
state._timers.pop(peer_id, None)
|
||||
teardown = use_cases.expire_session(peer_id)
|
||||
if teardown is None:
|
||||
return
|
||||
if rt["debug"]:
|
||||
logger.debug(
|
||||
"(PROXY) session timeout peer=%s client=%s:%s",
|
||||
int_id(peer_id),
|
||||
teardown.client.host,
|
||||
teardown.client.port,
|
||||
)
|
||||
if rt["client_info"] and peer_id != b"\xff\xff\xff\xff":
|
||||
logger.info(
|
||||
"(PROXY) Client: ID:%s IP:%s Port:%s Removed.",
|
||||
str(int_id(peer_id)).rjust(9),
|
||||
teardown.client.host.rjust(15),
|
||||
teardown.client.port,
|
||||
)
|
||||
apply_session_teardown(
|
||||
teardown,
|
||||
master_sink=master_sink,
|
||||
client_sender=state.client_sender,
|
||||
peer_registry=peer_registry,
|
||||
)
|
||||
|
||||
def _handle_prbl(data: bytes, addr: tuple[str, int]) -> None:
|
||||
expire = use_cases.block_ip_from_prbl(data, addr[0])
|
||||
if state._runtime["client_info"]:
|
||||
logger.info("(PROXY) Add to blacklist: host %s expire %s", addr[0], expire)
|
||||
|
||||
fanin._on_attached = _on_client_attached # noqa: SLF001
|
||||
fanin_proto, udp_port = listen_proxy_fanin(
|
||||
reactor,
|
||||
runtime["listen_ip"],
|
||||
runtime["listen_port"],
|
||||
use_cases,
|
||||
master_sink,
|
||||
debug=runtime["debug"],
|
||||
logger=logger,
|
||||
protocol=fanin,
|
||||
)
|
||||
state.udp_port = udp_port
|
||||
state.client_sender = FanInClientSender(fanin_proto.transport)
|
||||
target_proto.transport = ProxyReplyTransport(fanin_proto.transport, prbl_handler=_handle_prbl)
|
||||
|
||||
logger.info(
|
||||
"(PROXY) Hotspot fan-in on %s:%s → inject %s (MAX_PEERS=%s, TIMEOUT=%ss)",
|
||||
runtime["listen_ip"] or "*",
|
||||
runtime["listen_port"],
|
||||
target,
|
||||
runtime["max_peers"],
|
||||
runtime["timeout"],
|
||||
)
|
||||
return state
|
||||
@ -0,0 +1,95 @@
|
||||
"""UDP fan-in: hotspot LISTEN_PORT with in-process inject (Phase 3)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from collections.abc import Callable
|
||||
from typing import Any, Protocol
|
||||
|
||||
from twisted.internet.protocol import DatagramProtocol
|
||||
|
||||
from adn_server.application.ports import ProxyMasterSink
|
||||
from adn_server.application.proxy import ProxyUseCases, peer_id_from_packet
|
||||
from adn_server.domain.result import is_fail
|
||||
|
||||
_logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class _DatagramWriter(Protocol):
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
...
|
||||
|
||||
|
||||
class ProxyFanInProtocol(DatagramProtocol):
|
||||
"""UDP listener for hotspots; attaches sessions and injects into the target MASTER."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
proxy: ProxyUseCases,
|
||||
master_sink: ProxyMasterSink,
|
||||
*,
|
||||
debug: bool = False,
|
||||
logger: logging.Logger | None = None,
|
||||
on_attached: Callable[[bytes, str, int, bool], None] | None = None,
|
||||
) -> None:
|
||||
self._proxy = proxy
|
||||
self._master_sink = master_sink
|
||||
self.debug = debug
|
||||
self._log = logger or _logger
|
||||
self._on_attached = on_attached
|
||||
|
||||
def datagramReceived(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
host, port = addr
|
||||
if self._proxy.is_ip_blocked(host):
|
||||
if self.debug:
|
||||
self._log.debug("(PROXY) dropped packet from blacklisted IP %s:%s", host, port)
|
||||
return
|
||||
command = data[:4] if len(data) >= 4 else b""
|
||||
if self.debug:
|
||||
self._log.debug(
|
||||
"(PROXY) RX from %s:%s len=%d cmd=%r",
|
||||
host,
|
||||
port,
|
||||
len(data),
|
||||
command,
|
||||
)
|
||||
peer_id = peer_id_from_packet(data, from_master=False)
|
||||
if peer_id is None:
|
||||
if self.debug:
|
||||
self._log.debug("(PROXY) ignored packet with no peer_id from %s:%s", host, port)
|
||||
return
|
||||
new_session = self._proxy.resolve_client(peer_id) is None
|
||||
result = self._proxy.attach_client(peer_id, host, port)
|
||||
if is_fail(result):
|
||||
if self.debug:
|
||||
self._log.debug(
|
||||
"(PROXY) attach rejected peer=%s from %s:%s: %s",
|
||||
peer_id.hex(),
|
||||
host,
|
||||
port,
|
||||
result.error,
|
||||
)
|
||||
return
|
||||
if self._on_attached is not None:
|
||||
self._on_attached(peer_id, host, port, new_session)
|
||||
self._master_sink.inject(data, addr)
|
||||
|
||||
|
||||
def listen_proxy_fanin(
|
||||
reactor: Any,
|
||||
listen_ip: str,
|
||||
listen_port: int,
|
||||
proxy: ProxyUseCases,
|
||||
master_sink: ProxyMasterSink,
|
||||
*,
|
||||
debug: bool = False,
|
||||
logger: logging.Logger | None = None,
|
||||
protocol: ProxyFanInProtocol | None = None,
|
||||
) -> tuple[ProxyFanInProtocol, Any]:
|
||||
"""Bind LISTEN_PORT and return ``(protocol, udp_port)``."""
|
||||
fanin = protocol or ProxyFanInProtocol(proxy, master_sink, debug=debug, logger=logger)
|
||||
udp_port = reactor.listenUDP(listen_port, fanin, interface=listen_ip or "0.0.0.0")
|
||||
return fanin, udp_port
|
||||
|
||||
|
||||
__all__ = ["ProxyFanInProtocol", "listen_proxy_fanin"]
|
||||
@ -0,0 +1,31 @@
|
||||
"""Shared test fixtures."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
|
||||
def minimal_valid_config(**overrides: Any) -> dict[str, Any]:
|
||||
"""Minimal config dict that passes validate_config (integrated proxy required)."""
|
||||
config: dict[str, Any] = {
|
||||
"GLOBAL": {"SERVER_ID": 1},
|
||||
"REPORTS": {"REPORT": False},
|
||||
"SYSTEMS": {
|
||||
"HOTSPOT": {
|
||||
"MODE": "MASTER",
|
||||
"ENABLED": True,
|
||||
"MAX_PEERS": 8,
|
||||
}
|
||||
},
|
||||
"PROXY": {
|
||||
"LISTEN_PORT": 62031,
|
||||
"TARGET_SYSTEM": "HOTSPOT",
|
||||
"TIMEOUT": 30,
|
||||
},
|
||||
}
|
||||
for key, value in overrides.items():
|
||||
if isinstance(value, dict) and isinstance(config.get(key), dict):
|
||||
config[key] = {**config[key], **value}
|
||||
else:
|
||||
config[key] = value
|
||||
return config
|
||||
@ -0,0 +1,19 @@
|
||||
"""Regression: configs that must keep loading after proxy integration."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from adn_server.infrastructure import YamlConfigLoader
|
||||
|
||||
|
||||
def test_adn_parrot_yaml_loads_without_proxy_block() -> None:
|
||||
loader = YamlConfigLoader("/opt/new-adn-server")
|
||||
config = loader.load("/opt/new-adn-server/adn-parrot.yaml")
|
||||
assert "PROXY" not in config or config.get("PROXY") == {}
|
||||
assert config["SYSTEMS"]["PARROT"]["MODE"] == "PEER"
|
||||
|
||||
|
||||
def test_adn_server_yaml_loads_with_proxy() -> None:
|
||||
loader = YamlConfigLoader("/opt/new-adn-server")
|
||||
config = loader.load("/opt/new-adn-server/adn-server.yaml")
|
||||
assert config["PROXY"]["TARGET_SYSTEM"] == "SYSTEM"
|
||||
assert config["PROXY"]["LISTEN_PORT"] == 62031
|
||||
@ -0,0 +1,118 @@
|
||||
"""PROXY configuration validation (Phase 3)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
import pytest
|
||||
|
||||
from adn_server.application.proxy.deployment import (
|
||||
is_proxy_inject_only,
|
||||
normalize_proxy_target,
|
||||
)
|
||||
from adn_server.domain.errors import ConfigError
|
||||
from adn_server.infrastructure.config_normalizer import expand_generator
|
||||
from adn_server.infrastructure.config_validator import validate_config
|
||||
from adn_server.infrastructure.proxy.config import apply_proxy_env_overrides
|
||||
|
||||
from tests.conftest import minimal_valid_config
|
||||
|
||||
|
||||
def test_proxy_section_required_when_master_present() -> None:
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config({
|
||||
"GLOBAL": {"SERVER_ID": 1},
|
||||
"SYSTEMS": {"HOTSPOT": {"MODE": "MASTER", "ENABLED": True, "MAX_PEERS": 1}},
|
||||
})
|
||||
assert "PROXY" in str(exc.value)
|
||||
|
||||
|
||||
def test_parrot_config_without_proxy_is_valid() -> None:
|
||||
validate_config({
|
||||
"GLOBAL": {"SERVER_ID": 9990},
|
||||
"SYSTEMS": {"PARROT": {"MODE": "PEER", "ENABLED": True, "PORT": 54915}},
|
||||
})
|
||||
|
||||
|
||||
def test_proxy_enabled_key_rejected() -> None:
|
||||
config = minimal_valid_config()
|
||||
config["PROXY"]["ENABLED"] = False
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "PROXY.ENABLED" in str(exc.value)
|
||||
|
||||
|
||||
def test_proxy_enabled_requires_listen_port_and_target() -> None:
|
||||
config = minimal_valid_config()
|
||||
del config["PROXY"]["TARGET_SYSTEM"]
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "TARGET_SYSTEM" in str(exc.value)
|
||||
|
||||
|
||||
def test_proxy_target_rejects_port_and_generator() -> None:
|
||||
config = minimal_valid_config()
|
||||
config["SYSTEMS"]["HOTSPOT"]["PORT"] = 56400
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "PORT" in str(exc.value)
|
||||
|
||||
config = minimal_valid_config()
|
||||
config["SYSTEMS"]["HOTSPOT"]["GENERATOR"] = 102
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "GENERATOR" in str(exc.value)
|
||||
|
||||
|
||||
def test_proxy_rejects_deprecated_keys() -> None:
|
||||
config = minimal_valid_config()
|
||||
config["PROXY"]["DISPATCH"] = True
|
||||
with pytest.raises(ConfigError) as exc:
|
||||
validate_config(config)
|
||||
assert "DISPATCH" in str(exc.value)
|
||||
|
||||
|
||||
def test_valid_minimal_proxy_config() -> None:
|
||||
validate_config(minimal_valid_config())
|
||||
|
||||
|
||||
def test_normalize_proxy_target_strips_bind_fields() -> None:
|
||||
config = minimal_valid_config()
|
||||
config["SYSTEMS"]["HOTSPOT"]["PORT"] = 56400
|
||||
config["SYSTEMS"]["HOTSPOT"]["IP"] = "127.0.0.1"
|
||||
config["SYSTEMS"]["HOTSPOT"]["GENERATOR"] = 102
|
||||
normalize_proxy_target(config)
|
||||
target = config["SYSTEMS"]["HOTSPOT"]
|
||||
assert "PORT" not in target
|
||||
assert "IP" not in target
|
||||
assert target["_REPORT_BASE_PORT"] == 56400
|
||||
assert is_proxy_inject_only(config, "HOTSPOT")
|
||||
|
||||
|
||||
def test_apply_proxy_env_overrides(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
config = minimal_valid_config()
|
||||
monkeypatch.setenv("ADN_PROXY_DEBUG", "1")
|
||||
monkeypatch.setenv("ADN_PROXY_LISTENPORT", "63000")
|
||||
monkeypatch.setenv("ADN_PROXY_IPV6", "1")
|
||||
apply_proxy_env_overrides(config)
|
||||
assert config["PROXY"]["DEBUG"] is True
|
||||
assert config["PROXY"]["LISTEN_PORT"] == 63000
|
||||
assert config["PROXY"]["LISTEN_IP"] == "::"
|
||||
|
||||
|
||||
def test_expand_generator_skips_proxy_target() -> None:
|
||||
config = minimal_valid_config(
|
||||
PROXY={"LISTEN_PORT": 62031, "TARGET_SYSTEM": "HOTSPOT"},
|
||||
SYSTEMS={
|
||||
"HOTSPOT": {
|
||||
"MODE": "MASTER",
|
||||
"ENABLED": True,
|
||||
"PORT": 56400,
|
||||
"GENERATOR": 4,
|
||||
"MAX_PEERS": 8,
|
||||
}
|
||||
},
|
||||
)
|
||||
expand_generator(config, logging.getLogger("test"))
|
||||
assert "HOTSPOT-0" not in config["SYSTEMS"]
|
||||
assert "HOTSPOT" in config["SYSTEMS"]
|
||||
@ -0,0 +1,87 @@
|
||||
"""Proxy hot-reload keeps LISTEN_PORT and active sessions."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from adn_server.application.proxy import ProxyUseCases
|
||||
from adn_server.domain.proxy import ClientEndpoint, ClientSlot
|
||||
from adn_server.infrastructure.proxy.ip_blacklist import InMemoryProxyIpBlacklist
|
||||
from adn_server.infrastructure.proxy.rpto_queue import InMemoryPendingRptoQueue
|
||||
from adn_server.infrastructure.proxy.runtime import ProxyServiceState, apply_proxy_config_reload
|
||||
from adn_server.infrastructure.proxy.slot_store import InMemoryProxySlotStore
|
||||
from adn_server.infrastructure.proxy.udp_fanin import ProxyFanInProtocol
|
||||
|
||||
|
||||
def _minimal_proxy_state() -> ProxyServiceState:
|
||||
use_cases = ProxyUseCases(
|
||||
InMemoryProxySlotStore(),
|
||||
InMemoryPendingRptoQueue(),
|
||||
max_peers=10,
|
||||
black_list=(),
|
||||
ip_blacklist=InMemoryProxyIpBlacklist(),
|
||||
)
|
||||
peer = b"\x00\x07\x06\xf5" # 730039101
|
||||
use_cases._slots.bind( # noqa: SLF001
|
||||
ClientSlot(
|
||||
peer_id=peer,
|
||||
client=ClientEndpoint(host="10.0.0.1", port=62031),
|
||||
report_slot=0,
|
||||
)
|
||||
)
|
||||
fanin = ProxyFanInProtocol(use_cases, master_sink=None) # type: ignore[arg-type]
|
||||
return ProxyServiceState(
|
||||
target_system="SYSTEM",
|
||||
use_cases=use_cases,
|
||||
master_sink=None, # type: ignore[arg-type]
|
||||
client_sender=None, # type: ignore[arg-type]
|
||||
fanin=fanin,
|
||||
udp_port=object(),
|
||||
listen_port=62031,
|
||||
listen_ip="",
|
||||
_runtime={
|
||||
"listen_port": 62031,
|
||||
"listen_ip": "",
|
||||
"target_system": "SYSTEM",
|
||||
"timeout": 30.0,
|
||||
"debug": False,
|
||||
"client_info": True,
|
||||
"black_list": (),
|
||||
"ip_black_list": {},
|
||||
"max_peers": 10,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def test_apply_proxy_config_reload_keeps_sessions_and_updates_timeout() -> None:
|
||||
state = _minimal_proxy_state()
|
||||
config = {
|
||||
"PROXY": {
|
||||
"LISTEN_PORT": 62031,
|
||||
"TARGET_SYSTEM": "SYSTEM",
|
||||
"TIMEOUT": 45,
|
||||
"DEBUG": True,
|
||||
},
|
||||
"SYSTEMS": {"SYSTEM": {"MAX_PEERS": 20}},
|
||||
}
|
||||
log = logging.getLogger("test.proxy.reload")
|
||||
|
||||
apply_proxy_config_reload(state, config, logger=log)
|
||||
|
||||
assert len(state.use_cases.list_slots()) == 1
|
||||
assert state.udp_port is not None
|
||||
assert state._runtime["timeout"] == 45.0
|
||||
assert state.fanin.debug is True
|
||||
assert state.use_cases._max_peers == 20 # noqa: SLF001
|
||||
|
||||
|
||||
def test_proxy_use_cases_apply_runtime_settings() -> None:
|
||||
uc = ProxyUseCases(
|
||||
InMemoryProxySlotStore(),
|
||||
InMemoryPendingRptoQueue(),
|
||||
max_peers=5,
|
||||
black_list=(1001,),
|
||||
)
|
||||
uc.apply_runtime_settings(max_peers=12, black_list=(2002, 3003))
|
||||
assert uc._max_peers == 12 # noqa: SLF001
|
||||
assert uc._black_list == frozenset({2002, 3003}) # noqa: SLF001
|
||||
@ -0,0 +1,91 @@
|
||||
"""Live UDP smoke test for integrated proxy (isolated port; no production restart)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import socket
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
from twisted.internet import reactor
|
||||
|
||||
from adn_server.application.proxy import ProxyUseCases
|
||||
from adn_server.domain.value_objects import bytes_4
|
||||
from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config
|
||||
from adn_server.infrastructure.hbp_constants import RPTACK, RPTL
|
||||
from adn_server.infrastructure.proxy import (
|
||||
InMemoryPendingRptoQueue,
|
||||
InMemoryProxySlotStore,
|
||||
InProcessHbpSink,
|
||||
ProxyFanInProtocol,
|
||||
ProxyReplyTransport,
|
||||
)
|
||||
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
|
||||
|
||||
_SMOKE_PORT = 62032
|
||||
_PEER = bytes_4(1234567)
|
||||
|
||||
|
||||
class _AclRouter:
|
||||
def acl_check(self, peer_id: bytes, acl: object) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def _build_fanin() -> ProxyFanInProtocol:
|
||||
config = {
|
||||
"GLOBAL": {"PING_TIME": 10, "MAX_MISSED": 3, "USE_ACL": False},
|
||||
"SYSTEMS": {
|
||||
"HOTSPOT": {
|
||||
"MODE": "MASTER",
|
||||
"ENABLED": True,
|
||||
"MAX_PEERS": 8,
|
||||
"OPTIONS": "TS2=9990;",
|
||||
}
|
||||
},
|
||||
}
|
||||
ensure_system_runtime_config(config)
|
||||
hbp = HBPProtocol("HOTSPOT", config)
|
||||
hbp._router = _AclRouter() # type: ignore[assignment]
|
||||
proxy = ProxyUseCases(
|
||||
InMemoryProxySlotStore(),
|
||||
InMemoryPendingRptoQueue(),
|
||||
max_peers=8,
|
||||
)
|
||||
sink = InProcessHbpSink(hbp)
|
||||
fanin = ProxyFanInProtocol(proxy, sink)
|
||||
return fanin, hbp, sink
|
||||
|
||||
|
||||
@pytest.mark.smoke
|
||||
def test_live_udp_rptl_rptack_on_isolated_port() -> None:
|
||||
"""RPTL in → inject → RPTACK out on 127.0.0.1:62032 (does not use production 62031)."""
|
||||
fanin, hbp, _ = _build_fanin()
|
||||
result: dict[str, bytes | None] = {"reply": None, "error": None}
|
||||
|
||||
def _run_client() -> None:
|
||||
try:
|
||||
client = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||
client.settimeout(2.0)
|
||||
client.bind(("127.0.0.1", 0))
|
||||
client.sendto(RPTL + _PEER, ("127.0.0.1", _SMOKE_PORT))
|
||||
data, _ = client.recvfrom(4096)
|
||||
result["reply"] = data
|
||||
client.close()
|
||||
except OSError as exc:
|
||||
result["error"] = str(exc).encode()
|
||||
finally:
|
||||
reactor.callFromThread(reactor.stop)
|
||||
|
||||
listener = reactor.listenUDP(_SMOKE_PORT, fanin, interface="127.0.0.1")
|
||||
hbp.transport = ProxyReplyTransport(fanin.transport)
|
||||
reactor.callWhenRunning(
|
||||
lambda: threading.Thread(target=_run_client, daemon=True).start()
|
||||
)
|
||||
reactor.callLater(5.0, reactor.stop)
|
||||
reactor.run()
|
||||
|
||||
listener.stopListening()
|
||||
if result["error"]:
|
||||
pytest.fail(result["error"].decode())
|
||||
reply = result["reply"]
|
||||
assert reply is not None, "no UDP reply (timeout?)"
|
||||
assert reply.startswith(RPTACK), f"expected RPTACK, got {reply[:8]!r}"
|
||||
@ -0,0 +1,35 @@
|
||||
"""Session teardown executor (legacy reaper wire parity)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
from adn_server.application.proxy.session_teardown import CLIENT_TEARDOWN_REPEAT, client_teardown_packet, master_teardown_packet
|
||||
from adn_server.domain.proxy import ClientEndpoint, SessionTeardown
|
||||
from adn_server.infrastructure.proxy.session_executor import apply_session_teardown
|
||||
from adn_server.domain.value_objects import bytes_4
|
||||
|
||||
_PEER = bytes_4(1234567)
|
||||
_CLIENT = ClientEndpoint(host="10.0.0.8", port=5000)
|
||||
|
||||
|
||||
def test_apply_session_teardown_sends_rptcl_and_mstcl() -> None:
|
||||
master = MagicMock()
|
||||
client = MagicMock()
|
||||
registry = MagicMock()
|
||||
teardown = SessionTeardown(peer_id=_PEER, client=_CLIENT)
|
||||
apply_session_teardown(
|
||||
teardown,
|
||||
master_sink=master,
|
||||
client_sender=client,
|
||||
peer_registry=registry,
|
||||
)
|
||||
master.inject.assert_called_once_with(
|
||||
master_teardown_packet(_PEER),
|
||||
("10.0.0.8", 5000),
|
||||
)
|
||||
assert client.send_to_client.call_count == CLIENT_TEARDOWN_REPEAT
|
||||
for call in client.send_to_client.call_args_list:
|
||||
assert call.args[0] == client_teardown_packet()
|
||||
assert call.args[1] == _CLIENT
|
||||
registry.remove_peer.assert_called_once_with(_PEER)
|
||||
@ -0,0 +1,109 @@
|
||||
"""UDP fan-in protocol tests (no live reactor)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
from adn_server.application.proxy import ProxyUseCases
|
||||
from adn_server.domain.value_objects import bytes_4
|
||||
from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config
|
||||
from adn_server.infrastructure.hbp_constants import RPTACK, RPTL
|
||||
from adn_server.infrastructure.proxy import (
|
||||
InMemoryPendingRptoQueue,
|
||||
InMemoryProxySlotStore,
|
||||
InProcessHbpSink,
|
||||
ProxyFanInProtocol,
|
||||
ProxyReplyTransport,
|
||||
)
|
||||
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
|
||||
|
||||
_PEER = bytes_4(1234567)
|
||||
_CLIENT_ADDR = ("192.168.1.50", 62031)
|
||||
|
||||
|
||||
class _RecordingTransport:
|
||||
def __init__(self) -> None:
|
||||
self.sent: list[tuple[bytes, tuple[str, int]]] = []
|
||||
|
||||
def write(self, data: bytes, addr: tuple[str, int]) -> None:
|
||||
self.sent.append((data, addr))
|
||||
|
||||
|
||||
class _AclRouter:
|
||||
def acl_check(self, peer_id: bytes, acl: object) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def _hotspot_master_config() -> dict:
|
||||
config = {
|
||||
"GLOBAL": {"PING_TIME": 10, "MAX_MISSED": 3, "USE_ACL": False},
|
||||
"SYSTEMS": {
|
||||
"HOTSPOT": {
|
||||
"MODE": "MASTER",
|
||||
"ENABLED": True,
|
||||
"MAX_PEERS": 8,
|
||||
"OPTIONS": "TS2=9990;",
|
||||
}
|
||||
},
|
||||
}
|
||||
ensure_system_runtime_config(config)
|
||||
return config
|
||||
|
||||
|
||||
def _fanin_stack(*, max_peers: int = 8) -> tuple[ProxyFanInProtocol, ProxyUseCases, InProcessHbpSink, _RecordingTransport, MagicMock]:
|
||||
transport = _RecordingTransport()
|
||||
config = _hotspot_master_config()
|
||||
hbp = HBPProtocol("HOTSPOT", config)
|
||||
hbp._router = _AclRouter() # type: ignore[assignment]
|
||||
hbp.transport = ProxyReplyTransport(transport)
|
||||
sink = InProcessHbpSink(hbp)
|
||||
inject_spy = MagicMock(wraps=sink.inject)
|
||||
sink.inject = inject_spy # type: ignore[method-assign]
|
||||
proxy = ProxyUseCases(
|
||||
InMemoryProxySlotStore(),
|
||||
InMemoryPendingRptoQueue(),
|
||||
max_peers=max_peers,
|
||||
)
|
||||
protocol = ProxyFanInProtocol(proxy, sink)
|
||||
protocol.transport = transport # type: ignore[assignment]
|
||||
return protocol, proxy, sink, transport, inject_spy
|
||||
|
||||
|
||||
def test_client_packet_attaches_session_and_injects() -> None:
|
||||
protocol, proxy, _, _, inject_spy = _fanin_stack()
|
||||
packet = RPTL + _PEER
|
||||
protocol.datagramReceived(packet, _CLIENT_ADDR)
|
||||
inject_spy.assert_called_once_with(packet, _CLIENT_ADDR)
|
||||
slot = proxy.resolve_client(_PEER)
|
||||
assert slot is not None
|
||||
assert slot.host == _CLIENT_ADDR[0]
|
||||
assert slot.port == _CLIENT_ADDR[1]
|
||||
|
||||
|
||||
def test_existing_client_refreshes_endpoint_before_inject() -> None:
|
||||
protocol, proxy, _, _, inject_spy = _fanin_stack()
|
||||
packet = RPTL + _PEER
|
||||
protocol.datagramReceived(packet, _CLIENT_ADDR)
|
||||
inject_spy.reset_mock()
|
||||
new_addr = ("192.168.1.51", 62032)
|
||||
protocol.datagramReceived(packet, new_addr)
|
||||
inject_spy.assert_called_once_with(packet, new_addr)
|
||||
client = proxy.resolve_client(_PEER)
|
||||
assert client is not None
|
||||
assert client.host == new_addr[0]
|
||||
assert client.port == new_addr[1]
|
||||
|
||||
|
||||
def test_master_rptl_reply_sent_via_proxy_listen_socket() -> None:
|
||||
protocol, _, _, transport, _ = _fanin_stack()
|
||||
protocol.datagramReceived(RPTL + _PEER, _CLIENT_ADDR)
|
||||
assert len(transport.sent) == 1
|
||||
data, addr = transport.sent[0]
|
||||
assert data.startswith(RPTACK)
|
||||
assert addr == _CLIENT_ADDR
|
||||
|
||||
|
||||
def test_attach_rejection_skips_inject() -> None:
|
||||
protocol, _, _, _, inject_spy = _fanin_stack(max_peers=0)
|
||||
protocol.datagramReceived(RPTL + _PEER, _CLIENT_ADDR)
|
||||
inject_spy.assert_not_called()
|
||||
Loading…
Reference in new issue