Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
proxyserver.py397 linesDownload Raw Back to addons
1"""2This addon is responsible for starting/stopping the proxy server sockets/instances specified by the mode option.3"""4 5from __future__ import annotations6 7import asyncio8import collections9import ipaddress10import logging11from collections.abc import Iterable12from collections.abc import Iterator13from contextlib import contextmanager14from typing import Optional15 16from wsproto.frame_protocol import Opcode17 18from mitmproxy import command19from mitmproxy import ctx20from mitmproxy import exceptions21from mitmproxy import http22from mitmproxy import platform23from mitmproxy import tcp24from mitmproxy import udp25from mitmproxy import websocket26from mitmproxy.connection import Address27from mitmproxy.flow import Flow28from mitmproxy.proxy import events29from mitmproxy.proxy import mode_specs30from mitmproxy.proxy import server_hooks31from mitmproxy.proxy.layers.tcp import TcpMessageInjected32from mitmproxy.proxy.layers.udp import UdpMessageInjected33from mitmproxy.proxy.layers.websocket import WebSocketMessageInjected34from mitmproxy.proxy.mode_servers import ProxyConnectionHandler35from mitmproxy.proxy.mode_servers import ServerInstance36from mitmproxy.proxy.mode_servers import ServerManager37from mitmproxy.utils import asyncio_utils38from mitmproxy.utils import human39from mitmproxy.utils import signals40 41logger = logging.getLogger(__name__)42 43 44class Servers:45    def __init__(self, manager: ServerManager):46        self.changed = signals.AsyncSignal(lambda: None)47        self._instances: dict[mode_specs.ProxyMode, ServerInstance] = dict()48        self._lock = asyncio.Lock()49        self._manager = manager50 51    @property52    def is_updating(self) -> bool:53        return self._lock.locked()54 55    async def update(self, modes: Iterable[mode_specs.ProxyMode]) -> bool:56        all_ok = True57 58        async with self._lock:59            new_instances: dict[mode_specs.ProxyMode, ServerInstance] = {}60 61            start_tasks = []62            if ctx.options.server:63                # Create missing modes and keep existing ones.64                for spec in modes:65                    if spec in self._instances:66                        instance = self._instances[spec]67                    else:68                        instance = ServerInstance.make(spec, self._manager)69                        start_tasks.append(instance.start())70                    new_instances[spec] = instance71 72            # Shutdown modes that have been removed from the list.73            stop_tasks = [74                s.stop()75                for spec, s in self._instances.items()76                if spec not in new_instances77            ]78 79            if not start_tasks and not stop_tasks:80                return (81                    True  # nothing to do, so we don't need to trigger `self.changed`.82                )83 84            self._instances = new_instances85            # Notify listeners about the new not-yet-started servers.86            await self.changed.send()87 88            # We first need to free ports before starting new servers.89            for ret in await asyncio.gather(*stop_tasks, return_exceptions=True):90                if ret:91                    all_ok = False92                    logger.error(str(ret))93            for ret in await asyncio.gather(*start_tasks, return_exceptions=True):94                if ret:95                    all_ok = False96                    logger.error(str(ret))97 98        await self.changed.send()99        return all_ok100 101    def __len__(self) -> int:102        return len(self._instances)103 104    def __iter__(self) -> Iterator[ServerInstance]:105        return iter(self._instances.values())106 107    def __getitem__(self, mode: str | mode_specs.ProxyMode) -> ServerInstance:108        if isinstance(mode, str):109            mode = mode_specs.ProxyMode.parse(mode)110        return self._instances[mode]111 112 113class Proxyserver(ServerManager):114    """115    This addon runs the actual proxy server.116    """117 118    connections: dict[tuple | str, ProxyConnectionHandler]119    servers: Servers120 121    is_running: bool122    _connect_addr: Address | None = None123 124    def __init__(self):125        self.connections = {}126        self.servers = Servers(self)127        self.is_running = False128 129    def __repr__(self):130        return f"Proxyserver({len(self.connections)} active conns)"131 132    @command.command("proxyserver.active_connections")133    def active_connections(self) -> int:134        return len(self.connections)135 136    @contextmanager137    def register_connection(138        self, connection_id: tuple | str, handler: ProxyConnectionHandler139    ):140        self.connections[connection_id] = handler141        try:142            yield143        finally:144            del self.connections[connection_id]145 146    def load(self, loader):147        loader.add_option(148            "store_streamed_bodies",149            bool,150            False,151            "Store HTTP request and response bodies when streamed (see `stream_large_bodies`). "152            "This increases memory consumption, but makes it possible to inspect streamed bodies.",153        )154        loader.add_option(155            "connection_strategy",156            str,157            "eager",158            "Determine when server connections should be established. When set to lazy, mitmproxy "159            "tries to defer establishing an upstream connection as long as possible. This makes it possible to "160            "use server replay while being offline. When set to eager, mitmproxy can detect protocols with "161            "server-side greetings, as well as accurately mirror TLS ALPN negotiation.",162            choices=("eager", "lazy"),163        )164        loader.add_option(165            "stream_large_bodies",166            Optional[str],167            None,168            """169            Stream data to the client if request or response body exceeds the given170            threshold. If streamed, the body will not be stored in any way,171            and such responses cannot be modified. Understands k/m/g172            suffixes, i.e. 3m for 3 megabytes. To store streamed bodies, see `store_streamed_bodies`.173            """,174        )175        loader.add_option(176            "body_size_limit",177            Optional[str],178            None,179            """180            Byte size limit of HTTP request and response bodies. Understands181            k/m/g suffixes, i.e. 3m for 3 megabytes.182            """,183        )184        loader.add_option(185            "keep_host_header",186            bool,187            False,188            """189            Reverse Proxy: Keep the original host header instead of rewriting it190            to the reverse proxy target.191            """,192        )193        loader.add_option(194            "proxy_debug",195            bool,196            False,197            "Enable debug logs in the proxy core.",198        )199        loader.add_option(200            "normalize_outbound_headers",201            bool,202            True,203            """204            Normalize outgoing HTTP/2 header names, but emit a warning when doing so.205            HTTP/2 does not allow uppercase header names. This option makes sure that HTTP/2 headers set206            in custom scripts are lowercased before they are sent.207            """,208        )209        loader.add_option(210            "validate_inbound_headers",211            bool,212            True,213            """214            Make sure that incoming HTTP requests are not malformed.215            Disabling this option makes mitmproxy vulnerable to HTTP smuggling attacks.216            """,217        )218        loader.add_option(219            "connect_addr",220            Optional[str],221            None,222            """Set the local IP address that mitmproxy should use when connecting to upstream servers.""",223        )224 225    def running(self):226        self.is_running = True227 228    def configure(self, updated) -> None:229        if "stream_large_bodies" in updated:230            try:231                human.parse_size(ctx.options.stream_large_bodies)232            except ValueError:233                raise exceptions.OptionsError(234                    f"Invalid stream_large_bodies specification: "235                    f"{ctx.options.stream_large_bodies}"236                )237        if "body_size_limit" in updated:238            try:239                human.parse_size(ctx.options.body_size_limit)240            except ValueError:241                raise exceptions.OptionsError(242                    f"Invalid body_size_limit specification: "243                    f"{ctx.options.body_size_limit}"244                )245        if "connect_addr" in updated:246            try:247                if ctx.options.connect_addr:248                    self._connect_addr = (249                        str(ipaddress.ip_address(ctx.options.connect_addr)),250                        0,251                    )252                else:253                    self._connect_addr = None254            except ValueError:255                raise exceptions.OptionsError(256                    f"Invalid value for connect_addr: {ctx.options.connect_addr!r}. Specify a valid IP address."257                )258        if "mode" in updated or "server" in updated:259            # Make sure that all modes are syntactically valid...260            modes: list[mode_specs.ProxyMode] = []261            for mode in ctx.options.mode:262                try:263                    modes.append(mode_specs.ProxyMode.parse(mode))264                except ValueError as e:265                    raise exceptions.OptionsError(266                        f"Invalid proxy mode specification: {mode} ({e})"267                    )268 269            # ...and don't listen on the same address.270            listen_addrs = []271            for m in modes:272                if m.transport_protocol == "both":273                    protocols = ["tcp", "udp"]274                else:275                    protocols = [m.transport_protocol]276                host = m.listen_host(ctx.options.listen_host)277                port = m.listen_port(ctx.options.listen_port)278                if port is None:279                    continue280                for proto in protocols:281                    listen_addrs.append((host, port, proto))282            if len(set(listen_addrs)) != len(listen_addrs):283                (host, port, _) = collections.Counter(listen_addrs).most_common(1)[0][0]284                dup_addr = human.format_address((host or "0.0.0.0", port))285                raise exceptions.OptionsError(286                    f"Cannot spawn multiple servers on the same address: {dup_addr}"287                )288 289            if ctx.options.mode and not ctx.master.addons.get("nextlayer"):290                logger.warning("Warning: Running proxyserver without nextlayer addon!")291            if any(isinstance(m, mode_specs.TransparentMode) for m in modes):292                if platform.original_addr:293                    platform.init_transparent_mode()294                else:295                    raise exceptions.OptionsError(296                        "Transparent mode not supported on this platform."297                    )298 299            if self.is_running:300                asyncio_utils.create_task(301                    self.servers.update(modes),302                    name="update servers",303                    keep_ref=True,304                )305 306    async def setup_servers(self) -> bool:307        """Setup proxy servers. This may take an indefinite amount of time to complete (e.g. on permission prompts)."""308        return await self.servers.update(309            [mode_specs.ProxyMode.parse(m) for m in ctx.options.mode]310        )311 312    def listen_addrs(self) -> list[Address]:313        return [addr for server in self.servers for addr in server.listen_addrs]314 315    def inject_event(self, event: events.MessageInjected):316        connection_id: str | tuple317        if event.flow.client_conn.transport_protocol != "udp":318            connection_id = event.flow.client_conn.id319        else:  # pragma: no cover320            # temporary workaround: for UDP we don't have persistent client IDs yet.321            connection_id = (322                event.flow.client_conn.peername,323                event.flow.client_conn.sockname,324            )325        if connection_id not in self.connections:326            raise ValueError("Flow is not from a live connection.")327 328        asyncio_utils.create_task(329            self.connections[connection_id].server_event(event),330            name=f"inject_event",331            keep_ref=True,332            client=event.flow.client_conn.peername,333        )334 335    @command.command("inject.websocket")336    def inject_websocket(337        self, flow: Flow, to_client: bool, message: bytes, is_text: bool = True338    ):339        if not isinstance(flow, http.HTTPFlow) or not flow.websocket:340            logger.warning("Cannot inject WebSocket messages into non-WebSocket flows.")341            return342 343        msg = websocket.WebSocketMessage(344            Opcode.TEXT if is_text else Opcode.BINARY, not to_client, message345        )346        event = WebSocketMessageInjected(flow, msg)347        try:348            self.inject_event(event)349        except ValueError as e:350            logger.warning(str(e))351 352    @command.command("inject.tcp")353    def inject_tcp(self, flow: Flow, to_client: bool, message: bytes):354        if not isinstance(flow, tcp.TCPFlow):355            logger.warning("Cannot inject TCP messages into non-TCP flows.")356            return357 358        event = TcpMessageInjected(flow, tcp.TCPMessage(not to_client, message))359        try:360            self.inject_event(event)361        except ValueError as e:362            logger.warning(str(e))363 364    @command.command("inject.udp")365    def inject_udp(self, flow: Flow, to_client: bool, message: bytes):366        if not isinstance(flow, udp.UDPFlow):367            logger.warning("Cannot inject UDP messages into non-UDP flows.")368            return369 370        event = UdpMessageInjected(flow, udp.UDPMessage(not to_client, message))371        try:372            self.inject_event(event)373        except ValueError as e:374            logger.warning(str(e))375 376    def server_connect(self, data: server_hooks.ServerConnectionHookData):377        if data.server.sockname is None:378            data.server.sockname = self._connect_addr379 380        # Prevent mitmproxy from recursively connecting to itself.381        assert data.server.address382        connect_host, connect_port, *_ = data.server.address383 384        for server in self.servers:385            for listen_host, listen_port, *_ in server.listen_addrs:386                self_connect = (387                    connect_port == listen_port388                    and connect_host in ("localhost", "127.0.0.1", "::1", listen_host)389                    and server.mode.transport_protocol == data.server.transport_protocol390                )391                if self_connect:392                    data.server.error = (393                        "Request destination unknown. "394                        "Unable to figure out where this request should be forwarded to."395                    )396                    return397 
codekingpro/portable-devtools · Team Ai