codekingpro/portable-devtools
115k
1"""2Base class for protocol layers.3"""4 5import collections6import textwrap7from abc import abstractmethod8from collections.abc import Callable9from collections.abc import Generator10from dataclasses import dataclass11from logging import DEBUG12from typing import Any13from typing import ClassVar14from typing import NamedTuple15from typing import TypeVar16 17from mitmproxy.connection import Connection18from mitmproxy.proxy import commands19from mitmproxy.proxy import events20from mitmproxy.proxy.commands import Command21from mitmproxy.proxy.commands import StartHook22from mitmproxy.proxy.context import Context23 24T = TypeVar("T")25CommandGenerator = Generator[Command, Any, T]26"""27A function annotated with CommandGenerator[bool] may yield commands and ultimately return a boolean value.28"""29 30 31MAX_LOG_STATEMENT_SIZE = 204832"""Maximum size of individual log statements before they will be truncated."""33 34 35class Paused(NamedTuple):36 """37 State of a layer that's paused because it is waiting for a command reply.38 """39 40 command: commands.Command41 generator: CommandGenerator42 43 44class Layer:45 """46 The base class for all protocol layers.47 48 Layers interface with their child layer(s) by calling .handle_event(event),49 which returns a list (more precisely: a generator) of commands.50 Most layers do not implement .directly, but instead implement ._handle_event, which51 is called by the default implementation of .handle_event.52 The default implementation of .handle_event allows layers to emulate blocking code:53 When ._handle_event yields a command that has its blocking attribute set to True, .handle_event pauses54 the execution of ._handle_event and waits until it is called with the corresponding CommandCompleted event.55 All events encountered in the meantime are buffered and replayed after execution is resumed.56 57 The result is code that looks like blocking code, but is not blocking:58 59 def _handle_event(self, event):60 err = yield OpenConnection(server) # execution continues here after a connection has been established.61 62 Technically this is very similar to how coroutines are implemented.63 """64 65 __last_debug_message: ClassVar[str] = ""66 context: Context67 _paused: Paused | None68 """69 If execution is currently paused, this attribute stores the paused coroutine70 and the command for which we are expecting a reply.71 """72 _paused_event_queue: collections.deque[events.Event]73 """74 All events that have occurred since execution was paused.75 These will be replayed to ._child_layer once we resume.76 """77 debug: str | None = None78 """79 Enable debug logging by assigning a prefix string for log messages.80 Different amounts of whitespace for different layers work well.81 """82 83 def __init__(self, context: Context) -> None:84 self.context = context85 self.context.layers.append(self)86 self._paused = None87 self._paused_event_queue = collections.deque()88 89 show_debug_output = getattr(context.options, "proxy_debug", False)90 if show_debug_output: # pragma: no cover91 self.debug = " " * len(context.layers)92 93 def __repr__(self):94 statefun = getattr(self, "state", self._handle_event)95 state = getattr(statefun, "__name__", "")96 state = state.replace("state_", "")97 if state == "_handle_event":98 state = ""99 else:100 state = f"state: {state}"101 return f"{type(self).__name__}({state})"102 103 def __debug(self, message):104 """yield a Log command indicating what message is passing through this layer."""105 if len(message) > MAX_LOG_STATEMENT_SIZE:106 message = message[:MAX_LOG_STATEMENT_SIZE] + "…"107 if Layer.__last_debug_message == message:108 message = message.split("\n", 1)[0].strip()109 if len(message) > 256:110 message = message[:256] + "…"111 else:112 Layer.__last_debug_message = message113 assert self.debug is not None114 return commands.Log(textwrap.indent(message, self.debug), DEBUG)115 116 @property117 def stack_pos(self) -> str:118 """repr() for this layer and all its parent layers, only useful for debugging."""119 try:120 idx = self.context.layers.index(self)121 except ValueError:122 return repr(self)123 else:124 return " >> ".join(repr(x) for x in self.context.layers[: idx + 1])125 126 @abstractmethod127 def _handle_event(self, event: events.Event) -> CommandGenerator[None]:128 """Handle a proxy server event"""129 yield from () # pragma: no cover130 131 def handle_event(self, event: events.Event) -> CommandGenerator[None]:132 if self._paused:133 # did we just receive the reply we were waiting for?134 pause_finished = (135 isinstance(event, events.CommandCompleted)136 and event.command is self._paused.command137 )138 if self.debug is not None:139 yield self.__debug(f"{'>>' if pause_finished else '>!'} {event}")140 if pause_finished:141 assert isinstance(event, events.CommandCompleted)142 yield from self.__continue(event)143 else:144 self._paused_event_queue.append(event)145 else:146 if self.debug is not None:147 yield self.__debug(f">> {event}")148 command_generator = self._handle_event(event)149 send = None150 151 # inlined copy of __process to reduce call stack.152 # <✂✂✂>153 try:154 # Run ._handle_event to the next yield statement.155 # If you are not familiar with generators and their .send() method,156 # https://stackoverflow.com/a/12638313/934719 has a good explanation.157 command = command_generator.send(send)158 except StopIteration:159 return160 161 while True:162 if self.debug is not None:163 if not isinstance(command, commands.Log):164 yield self.__debug(f"<< {command}")165 if command.blocking is True:166 # We only want this layer to block, the outer layers should not block.167 # For example, take an HTTP/2 connection: If we intercept one particular request,168 # we don't want all other requests in the connection to be blocked a well.169 # We signal to outer layers that this command is already handled by assigning our layer to170 # `.blocking` here (upper layers explicitly check for `is True`).171 command.blocking = self172 self._paused = Paused(173 command,174 command_generator,175 )176 yield command177 return178 else:179 yield command180 try:181 command = next(command_generator)182 except StopIteration:183 return184 # </✂✂✂>185 186 def __process(self, command_generator: CommandGenerator, send=None):187 """188 Yield commands from a generator.189 If a command is blocking, execution is paused and this function returns without190 processing any further commands.191 """192 try:193 # Run ._handle_event to the next yield statement.194 # If you are not familiar with generators and their .send() method,195 # https://stackoverflow.com/a/12638313/934719 has a good explanation.196 command = command_generator.send(send)197 except StopIteration:198 return199 200 while True:201 if self.debug is not None:202 if not isinstance(command, commands.Log):203 yield self.__debug(f"<< {command}")204 if command.blocking is True:205 # We only want this layer to block, the outer layers should not block.206 # For example, take an HTTP/2 connection: If we intercept one particular request,207 # we don't want all other requests in the connection to be blocked a well.208 # We signal to outer layers that this command is already handled by assigning our layer to209 # `.blocking` here (upper layers explicitly check for `is True`).210 command.blocking = self211 self._paused = Paused(212 command,213 command_generator,214 )215 yield command216 return217 else:218 yield command219 try:220 command = next(command_generator)221 except StopIteration:222 return223 224 def __continue(self, event: events.CommandCompleted):225 """226 Continue processing events after being paused.227 The tricky part here is that events in the event queue may trigger commands which again pause the execution,228 so we may not be able to process the entire queue.229 """230 assert self._paused is not None231 command_generator = self._paused.generator232 self._paused = None233 yield from self.__process(command_generator, event.reply)234 235 while not self._paused and self._paused_event_queue:236 ev = self._paused_event_queue.popleft()237 if self.debug is not None:238 yield self.__debug(f"!> {ev}")239 command_generator = self._handle_event(ev)240 yield from self.__process(command_generator)241 242 243mevents = (244 events # alias here because autocomplete above should not have aliased version.245)246 247 248class NextLayer(Layer):249 layer: Layer | None250 """The next layer. To be set by an addon."""251 252 events: list[mevents.Event]253 """All events that happened before a decision was made."""254 255 _ask_on_start: bool256 257 def __init__(self, context: Context, ask_on_start: bool = False) -> None:258 super().__init__(context)259 self.context.layers.remove(self)260 self.layer = None261 self.events = []262 self._ask_on_start = ask_on_start263 self._handle: Callable[[mevents.Event], CommandGenerator[None]] | None = None264 265 def __repr__(self):266 return f"NextLayer:{self.layer!r}"267 268 def handle_event(self, event: mevents.Event):269 if self._handle is not None:270 yield from self._handle(event)271 else:272 yield from super().handle_event(event)273 274 def _handle_event(self, event: mevents.Event):275 self.events.append(event)276 277 # We receive new data. Let's find out if we can determine the next layer now?278 if self._ask_on_start and isinstance(event, events.Start):279 yield from self._ask()280 elif (281 isinstance(event, mevents.ConnectionClosed)282 and event.connection == self.context.client283 ):284 # If we have not determined the next protocol yet and the client already closes the connection,285 # we abort everything.286 yield commands.CloseConnection(self.context.client)287 elif isinstance(event, mevents.DataReceived):288 # For now, we only ask if we have received new data to reduce hook noise.289 yield from self._ask()290 291 def _ask(self):292 """293 Manually trigger a next_layer hook.294 The only use at the moment is to make sure that the top layer is initialized.295 """296 yield NextLayerHook(self)297 298 # Has an addon decided on the next layer yet?299 if self.layer:300 if self.debug:301 yield commands.Log(f"{self.debug}[nextlayer] {self.layer!r}", DEBUG)302 for e in self.events:303 yield from self.layer.handle_event(e)304 self.events.clear()305 306 # Why do we need three assignments here?307 # 1. When this function here is invoked we may have paused events. Those should be308 # forwarded to the sublayer right away, so we reassign ._handle_event.309 # 2. This layer is not needed anymore, so we directly reassign .handle_event.310 # 3. Some layers may however still have a reference to the old .handle_event.311 # ._handle is just an optimization to reduce the callstack in these cases.312 self.handle_event = self.layer.handle_event # type: ignore313 self._handle_event = self.layer.handle_event # type: ignore314 self._handle = self.layer.handle_event315 316 # Utility methods for whoever decides what the next layer is going to be.317 def data_client(self):318 return self._data(self.context.client)319 320 def data_server(self):321 return self._data(self.context.server)322 323 def _data(self, connection: Connection):324 data = (325 e.data326 for e in self.events327 if isinstance(e, mevents.DataReceived) and e.connection == connection328 )329 return b"".join(data)330 331 332@dataclass333class NextLayerHook(StartHook):334 """335 Network layers are being switched. You may change which layer will be used by setting data.layer.336 337 (by default, this is done by mitmproxy.addons.NextLayer)338 """339 340 data: NextLayer341 