codekingpro/portable-devtools
114k
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 