Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
layer.py341 linesDownload Raw Back to proxy
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 
codekingpro/portable-devtools · Team Ai