← pod_r4_v1g

httpx_3672

resolved RESOLVED PASS · 41 tool calls · 215 s · encode/httpx

Task input

Server connection handling.

* Add `HTTPParser.keep_alive`.
* Server... always read request to completion on keep alives.
* `HTTPParser.complete` -> `.reset`
* Close streams on server exit.
* Don't raise `KeyboardException` on server exit.

Tool calls (42)

#ToolArgumentsResult
1run_command
{
 "command": "grep -rn \"HTTPParser\" src/httpx/_parsers.py src/httpx/_server.py"
}
{"status": "ok", "stdout": "src/httpx/_parsers.py:5:__all__ = ['HTTPParser', 'Mode', 'ProtocolError']\nsrc/httpx/_parsers.py:72:class HTTPParser:\nsrc/httpx/_pa
full
{"status": "ok", "stdout": "src/httpx/_parsers.py:5:__all__ = ['HTTPParser', 'Mode', 'ProtocolError']\nsrc/httpx/_parsers.py:72:class HTTPParser:\nsrc/httpx/_parsers.py:76:    client = HTTPParser(writer, reader)\nsrc/httpx/_parsers.py:426:        return f'<HTTPParser [{detail}]>'\nsrc/httpx/_server.py:6:from ._parsers import HTTPParser\nsrc/httpx/_server.py:27:        self._parser = HTTPParser(stream, mode='SERVER')\n", "stderr": "", "exit_code": 0}
2run_command
{
 "command": "sed -n '72,150p' src/httpx/_parsers.py"
}
{"status": "ok", "stdout": "class HTTPParser:\n \"\"\"\n Usage...\n\n client = HTTPParser(writer, reader)\n client.send_method_line()\n client.send_headers()\n
full
{"status": "ok", "stdout": "class HTTPParser:\n    \"\"\"\n    Usage...\n\n    client = HTTPParser(writer, reader)\n    client.send_method_line()\n    client.send_headers()\n    client.send_body()\n    client.recv_status_line()\n    client.recv_headers()\n    client.recv_body()\n    client.complete()\n    client.close()\n    \"\"\"\n    def __init__(self, stream: Stream, mode: str) -> None:\n        self.stream = stream\n        self.parser = ReadAheadParser(stream)\n        self.mode = {'CLIENT': Mode.CLIENT, 'SERVER': Mode.SERVER}[mode]\n\n        # Track state...\n        if self.mode == Mode.CLIENT:\n            self.send_state: State = State.SEND_METHOD_LINE\n            self.recv_state: State = State.WAIT\n        else:\n            self.recv_state = State.RECV_METHOD_LINE\n            self.send_state = State.WAIT\n\n        # Track message framing...\n        self.send_content_length: int | None = 0\n        self.recv_content_length: int | None = 0\n        self.send_seen_length = 0\n        self.recv_seen_length = 0\n\n        # Track connection keep alive...\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\n\n        # Special states...\n        self.processing_1xx = False\n\n    def send_method_line(self, method: bytes, target: bytes, protocol: bytes) -> None:\n        \"\"\"\n        Send the initial request line:\n\n        >>> p.send_method_line(b'GET', b'/', b'HTTP/1.1')\n\n        Sending state will switch to SEND_HEADERS state.\n        \"\"\"\n        if self.send_state != State.SEND_METHOD_LINE:\n            msg = f\"Called 'send_method_line' in invalid state {self.send_state}\"\n            raise ProtocolError(msg)\n\n        # Send initial request line, eg. \"GET / HTTP/1.1\"\n        if protocol != b'HTTP/1.1':\n            raise ProtocolError(\"Sent unsupported protocol version\")\n        data = b\" \".join([method, target, protocol]) + b\"\\r\\n\"\n        self.stream.write(data)\n\n        self.send_state = State.SEN
3run_command
{
 "command": "grep -n \"def complete\" src/httpx/_parsers.py"
}
{"status": "ok", "stdout": "378: def complete(self):\n", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "378:    def complete(self):\n", "stderr": "", "exit_code": 0}
4run_command
{
 "command": "sed -n '378,385p' src/httpx/_parsers.py"
}
{"status": "ok", "stdout": " def complete(self):\n is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n is_keepalive = self.sen
full
{"status": "ok", "stdout": "    def complete(self):\n        is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n        is_keepalive = self.send_keep_alive and self.recv_keep_alive\n\n        if not (is_fully_complete and is_keepalive):\n            self.close()\n            return\n\n", "stderr": "", "exit_code": 0}
5run_command
{
 "command": "sed -n '378,400p' src/httpx/_parsers.py"
}
{"status": "ok", "stdout": " def complete(self):\n is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n is_keepalive = self.sen
full
{"status": "ok", "stdout": "    def complete(self):\n        is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n        is_keepalive = self.send_keep_alive and self.recv_keep_alive\n\n        if not (is_fully_complete and is_keepalive):\n            self.close()\n            return\n\n        if self.mode == Mode.CLIENT:\n            self.send_state = State.SEND_METHOD_LINE\n            self.recv_state = State.WAIT\n        else:\n            self.recv_state = State.RECV_METHOD_LINE\n            self.send_state = State.WAIT\n\n        self.send_content_length = 0\n        self.recv_content_length = 0\n        self.send_seen_length = 0\n        self.recv_seen_length = 0\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\n        self.processing_1xx = False\n\n", "stderr": "", "exit_code": 0}
6run_command
{
 "command": "sed -n '1,100p' src/httpx/_server.py"
}
{"status": "ok", "stdout": "import contextlib\nimport logging\nimport time\n\nfrom ._content import Text\nfrom ._parsers import HTTPParser\nfrom ._request impor
full
{"status": "ok", "stdout": "import contextlib\nimport logging\nimport time\n\nfrom ._content import Text\nfrom ._parsers import HTTPParser\nfrom ._request import Request\nfrom ._response import Response\nfrom ._network import NetworkBackend, sleep\nfrom ._streams import HTTPStream\n\n__all__ = [\n    \"serve_http\", \"run\"\n]\n\nlogger = logging.getLogger(\"httpx.server\")\n\n\nclass ConnectionClosed(Exception):\n    pass\n\n\nclass HTTPConnection:\n    def __init__(self, stream, endpoint):\n        self._stream = stream\n        self._endpoint = endpoint\n        self._parser = HTTPParser(stream, mode='SERVER')\n        self._keepalive_duration = 5.0\n        self._idle_expiry = time.monotonic() + self._keepalive_duration\n\n    # API entry points...\n    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\"Internal Server Error\", exc_info=True)\n                        content = Text(\"Internal Server Error\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n        except Exception:\n            logger.error(\"Internal Server Error\", exc_info=True)\n\n    def close(self):\n        self
7run_command
{
 "command": "sed -n '101,150p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " while(True):\n sleep(1)\n\n\n@contextlib.contextmanager\ndef serve_http(endpoint):\n def handler(stream):\n connection = HTTPConnec
full
{"status": "ok", "stdout": "        while(True):\n            sleep(1)\n\n\n@contextlib.contextmanager\ndef serve_http(endpoint):\n    def handler(stream):\n        connection = HTTPConnection(stream, endpoint)\n        connection.handle_requests()\n\n    logging.basicConfig(\n        format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n        datefmt=\"%Y-%m-%d %H:%M:%S\",\n        level=logging.DEBUG\n    )\n\n    backend = NetworkBackend()\n    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n        yield server\n\n\ndef run(app):\n    with serve_http(app) as server:\n        server.wait()\n", "stderr": "", "exit_code": 0}
8run_command
{
 "command": "grep -n \"def __init__\" src/httpx/_parsers.py"
}
{"status": "ok", "stdout": "86: def __init__(self, stream: Stream, mode: str) -> None:\n433: def __init__(self, stream: Stream) -> None:\n", "stderr": "", "exit
full
{"status": "ok", "stdout": "86:    def __init__(self, stream: Stream, mode: str) -> None:\n433:    def __init__(self, stream: Stream) -> None:\n", "stderr": "", "exit_code": 0}
9run_command
{
 "command": "sed -n '86,100p' src/httpx/_parsers.py"
}
{"status": "ok", "stdout": " def __init__(self, stream: Stream, mode: str) -> None:\n self.stream = stream\n self.parser = ReadAheadParser(stream)\n self.mode =
full
{"status": "ok", "stdout": "    def __init__(self, stream: Stream, mode: str) -> None:\n        self.stream = stream\n        self.parser = ReadAheadParser(stream)\n        self.mode = {'CLIENT': Mode.CLIENT, 'SERVER': Mode.SERVER}[mode]\n\n        # Track state...\n        if self.mode == Mode.CLIENT:\n            self.send_state: State = State.SEND_METHOD_LINE\n            self.recv_state: State = State.WAIT\n        else:\n            self.recv_state = State.RECV_METHOD_LINE\n            self.send_state = State.WAIT\n\n        # Track message framing...\n        self.send_content_length: int | None = 0\n", "stderr": "", "exit_code": 0}
10run_command
{
 "command": "sed -n '101,115p' src/httpx/_parsers.py"
}
{"status": "ok", "stdout": " self.recv_content_length: int | None = 0\n self.send_seen_length = 0\n self.recv_seen_length = 0\n\n # Track connection keep alive.
full
{"status": "ok", "stdout": "        self.recv_content_length: int | None = 0\n        self.send_seen_length = 0\n        self.recv_seen_length = 0\n\n        # Track connection keep alive...\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\n\n        # Special states...\n        self.processing_1xx = False\n\n    def send_method_line(self, method: bytes, target: bytes, protocol: bytes) -> None:\n        \"\"\"\n        Send the initial request line:\n\n", "stderr": "", "exit_code": 0}
11run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_parsers.py\")\ns = p.read_text()\nold = \"\"\"    def complete(self):\n        is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n        is_keepalive = self.send_keep_alive and self.recv_keep_alive\n\n        if not (is_fully_complete and is_keepalive):\n            self.close()\n            return\n\n        if self.mode == Mode.CLIENT:\n            self.send_state = State.SEND_METHOD_LINE\n            self.recv_state = State.WAIT\n        else:\n            self.recv_state = State.RECV_METHOD_LINE\n            self.send_state = State.WAIT\n\n        self.send_content_length = 0\n        self.recv_content_length = 0\n        self.send_seen_length = 0\n        self.recv_seen_length = 0\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\n        self.processing_1xx = False\"\"\"\nnew = \"\"\"    def reset(self):\n        is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n        is_keepalive = self.send_keep_alive and self.recv_keep_alive\n\n        if not (is_fully_complete and is_keepalive):\n            self.close()\n            return\n\n        if self.mode == Mode.CLIENT:\n            self.send_state = State.SEND_METHOD_LINE\n            self.recv_state = State.WAIT\n        else:\n            self.recv_state = State.RECV_METHOD_LINE\n            self.send_state = State.WAIT\n\n        self.send_content_length = 0\n        self.recv_content_length = 0\n        self.send_seen_length = 0\n        self.recv_seen_length = 0\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\n        self.processing_1xx = False\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
12run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_parsers.py\")\ns = p.read_text()\nold = \"\"\"        # Track connection keep alive...\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\"\"\"\nnew = \"\"\"        # Track connection keep alive...\n        self.send_keep_alive = True\n        self.recv_keep_alive = True\n        self.keep_alive = True\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
13run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_server.py\")\ns = p.read_text()\nold = \"\"\"    # Start it all over again...\n    def _complete(self):\n        self._parser.complete\n        self._idle_expiry = time.monotonic() + self._keepalive_duration\"\"\"\nnew = \"\"\"    # Start it all over again...\n    def _complete(self):\n        self._parser.reset()\n        self._idle_expiry = time.monotonic() + self._keepalive_duration\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
14run_command
{
 "command": "sed -n '20,60p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " pass\n\n\nclass HTTPConnection:\n def __init__(self, stream, endpoint):\n self._stream = stream\n self._endpoint = endpoint\n self.
full
{"status": "ok", "stdout": "    pass\n\n\nclass HTTPConnection:\n    def __init__(self, stream, endpoint):\n        self._stream = stream\n        self._endpoint = endpoint\n        self._parser = HTTPParser(stream, mode='SERVER')\n        self._keepalive_duration = 5.0\n        self._idle_expiry = time.monotonic() + self._keepalive_duration\n\n    # API entry points...\n    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\"Internal Server Error\", exc_info=True)\n                        content = Text(\"Internal Server Error\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n        except Exception:\n            logger.error(\"Internal Server Error\", exc_info=True)\n\n    def close(self):\n        self._parser.close()\n\n    # Receive the request...\n    def _recv_head(self) -> tuple[str, str, list[tuple[str, str]]]:\n        method, target, _ = self._parser.recv_method_line()\n", "stderr": "", "exit_code": 0}
15run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_server.py\")\ns = p.read_text()\nold = \"\"\"    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\\\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\\\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\\\"Internal Server Error\\\", exc_info=True)\n                        content = Text(\\\"Internal Server Error\\\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n        except Exception:\n            logger.error(\\\"Internal Server Error\\\", exc_info=True)\"\"\"\nnew = \"\"\"    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\\\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\\\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\\\"Internal Server Error\\\", exc_info=True)\n                        content = Text(\\\"Internal Server Error\\\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n                if self._parser.keep_alive:\n                    # Server... always read request to completion on keep alives.\n                    # This is not explicitly clear what it means, but usually it means\n                    # consuming the remaining body if the user didn't.\n                    # However, the Request context manager should handle it?\n                    # Let's check Request.__exit__.\n                    pass\n        except Exception:\n            logger.error(\\\"Internal Server Error\\\", exc_info=True)\"\"\"\n# Wait, I should check Request.__exit__ first.\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
16run_command
{
 "command": "grep -n \"class Request\" src/httpx/_request.py"
}
{"status": "ok", "stdout": "12:class Request:\n", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "12:class Request:\n", "stderr": "", "exit_code": 0}
17run_command
{
 "command": "sed -n '12,100p' src/httpx/_request.py"
}
{"status": "ok", "stdout": "class Request:\n def __init__(\n self,\n method: str,\n url: URL | str,\n headers: Headers | typing.Mapping[str, str] | None = None,
full
{"status": "ok", "stdout": "class Request:\n    def __init__(\n        self,\n        method: str,\n        url: URL | str,\n        headers: Headers | typing.Mapping[str, str] | None = None,\n        content: Content | Stream | bytes | None = None,\n    ):\n        self.method = method\n        self.url = URL(url)\n        self.headers = Headers(headers)\n        self.stream: Stream = ByteStream(b\"\")\n\n        # https://datatracker.ietf.org/doc/html/rfc2616#section-14.23\n        #\u00a0RFC 2616, Section 14.23, Host.\n        #\n        # A client MUST include a Host header field in all HTTP/1.1 request messages.\n        if \"Host\" not in self.headers:\n            self.headers = self.headers.copy_set(\"Host\", self.url.netloc)\n\n        if content is not None:\n            if isinstance(content, bytes):\n                self.stream = ByteStream(content)\n            elif isinstance(content, Stream):\n                self.stream = content\n            elif isinstance(content, Content):\n                ct = content.content_type()\n                self.stream = content.encode()\n                self.headers = self.headers.copy_set(\"Content-Type\", ct)\n            else:\n                raise TypeError(f'Expected `Content | Stream | bytes | None` got {type(content)}')\n\n            # https://datatracker.ietf.org/doc/html/rfc2616#section-4.3\n            #\u00a0RFC 2616, Section 4.3, Message Body.\n            #\n            # The presence of a message-body in a request is signaled by the\n            # inclusion of a Content-Length or Transfer-Encoding header field in\n            # the request's message-headers.\n            content_length: int | None = self.stream.size\n            if content_length is None:\n                self.headers = self.headers.copy_set(\"Transfer-Encoding\", \"chunked\")\n            elif content_length > 0:\n                self.headers = self.headers.copy_set(\"Content-Length\", str(content_length))\n\n        elif method in (\"
18run_command
{
 "command": "sed -n '110,150p' src/httpx/_server.py"
}
{"status": "ok", "stdout": "\n logging.basicConfig(\n format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n datefmt=\"%Y-%m-%d %H:%M:%S\",\n level=lo
full
{"status": "ok", "stdout": "\n    logging.basicConfig(\n        format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n        datefmt=\"%Y-%m-%d %H:%M:%S\",\n        level=logging.DEBUG\n    )\n\n    backend = NetworkBackend()\n    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n        yield server\n\n\ndef run(app):\n    with serve_http(app) as server:\n        server.wait()\n", "stderr": "", "exit_code": 0}
19run_command
{
 "command": "sed -n '101,110p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " while(True):\n sleep(1)\n\n\n@contextlib.contextmanager\ndef serve_http(endpoint):\n def handler(stream):\n connection = HTTPConnec
full
{"status": "ok", "stdout": "        while(True):\n            sleep(1)\n\n\n@contextlib.contextmanager\ndef serve_http(endpoint):\n    def handler(stream):\n        connection = HTTPConnection(stream, endpoint)\n        connection.handle_requests()\n\n", "stderr": "", "exit_code": 0}
20run_command
{
 "command": "sed -n '111,120p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " logging.basicConfig(\n format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n datefmt=\"%Y-%m-%d %H:%M:%S\",\n level=logg
full
{"status": "ok", "stdout": "    logging.basicConfig(\n        format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n        datefmt=\"%Y-%m-%d %H:%M:%S\",\n        level=logging.DEBUG\n    )\n\n    backend = NetworkBackend()\n    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n", "stderr": "", "exit_code": 0}
21-22 ×2run_command
{
 "command": "sed -n '121,130p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " yield server\n\n\ndef run(app):\n with serve_http(app) as server:\n server.wait()\n", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "        yield server\n\n\ndef run(app):\n    with serve_http(app) as server:\n        server.wait()\n", "stderr": "", "exit_code": 0}
23run_command
{
 "command": "sed -n '1,100p' src/httpx/_server.py"
}
{"status": "ok", "stdout": "import contextlib\nimport logging\nimport time\n\nfrom ._content import Text\nfrom ._parsers import HTTPParser\nfrom ._request impor
full
{"status": "ok", "stdout": "import contextlib\nimport logging\nimport time\n\nfrom ._content import Text\nfrom ._parsers import HTTPParser\nfrom ._request import Request\nfrom ._response import Response\nfrom ._network import NetworkBackend, sleep\nfrom ._streams import HTTPStream\n\n__all__ = [\n    \"serve_http\", \"run\"\n]\n\nlogger = logging.getLogger(\"httpx.server\")\n\n\nclass ConnectionClosed(Exception):\n    pass\n\n\nclass HTTPConnection:\n    def __init__(self, stream, endpoint):\n        self._stream = stream\n        self._endpoint = endpoint\n        self._parser = HTTPParser(stream, mode='SERVER')\n        self._keepalive_duration = 5.0\n        self._idle_expiry = time.monotonic() + self._keepalive_duration\n\n    # API entry points...\n    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\"Internal Server Error\", exc_info=True)\n                        content = Text(\"Internal Server Error\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n        except Exception:\n            logger.error(\"Internal Server Error\", exc_info=True)\n\n    def close(self):\n        self
24run_command
{
 "command": "sed -n '101,130p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " while(True):\n sleep(1)\n\n\n@contextlib.contextmanager\ndef serve_http(endpoint):\n def handler(stream):\n connection = HTTPConnec
full
{"status": "ok", "stdout": "        while(True):\n            sleep(1)\n\n\n@contextlib.contextmanager\ndef serve_http(endpoint):\n    def handler(stream):\n        connection = HTTPConnection(stream, endpoint)\n        connection.handle_requests()\n\n    logging.basicConfig(\n        format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n        datefmt=\"%Y-%m-%d %H:%M:%S\",\n        level=logging.DEBUG\n    )\n\n    backend = NetworkBackend()\n    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n        yield server\n\n\ndef run(app):\n    with serve_http(app) as server:\n        server.wait()\n", "stderr": "", "exit_code": 0}
25run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_server.py\")\ns = p.read_text()\nold = \"\"\"    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\\\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\\\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\\\"Internal Server Error\\\", exc_info=True)\n                        content = Text(\\\"Internal Server Error\\\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n        except Exception:\n            logger.error(\\\"Internal Server Error\\\", exc_info=True)\"\"\"\nnew = \"\"\"    def handle_requests(self):\n        try:\n            while not self._parser.is_closed():\n                method, url, headers = self._recv_head()\n                stream = HTTPStream(self._recv_body, self._complete)\n                # TODO: Handle endpoint exceptions\n                with Request(method, url, headers=headers, content=stream) as request:\n                    try:\n                        response = self._endpoint(request)\n                        status_line = f\\\"{request.method} {request.url.target} [{response.status_code} {response.reason_phrase}]\\\"\n                        logger.info(status_line)\n                    except Exception:\n                        logger.error(\\\"Internal Server Error\\\", exc_info=True)\n                        content = Text(\\\"Internal Server Error\\\")\n                        err = Response(code=500, content=content)\n                        self._send_head(err)\n                        self._send_body(err)\n                    else:\n                        self._send_head(response)\n                        self._send_body(response)\n                if self._parser.keep_alive:\n                    # Server... always read request to completion on keep alives.\n                    # The Request context manager calls self.close(), which calls self.stream.close().\n                    # For HTTPStream, close() should consume the remaining body.\n                    pass\n        except Exception:\n            logger.error(\\\"Internal Server Error\\\", exc_info=True)\"\"\"\n# I need to check HTTPStream.close()\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
26run_command
{
 "command": "grep -n \"class HTTPStream\" src/httpx/_streams.py"
}
{"status": "ok", "stdout": "94:class HTTPStream(Stream):\n", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "94:class HTTPStream(Stream):\n", "stderr": "", "exit_code": 0}
27run_command
{
 "command": "sed -n '94,150p' src/httpx/_streams.py"
}
{"status": "ok", "stdout": "class HTTPStream(Stream):\n def __init__(self, next_chunk, complete):\n self._next_chunk = next_chunk\n self._complete = complete\n
full
{"status": "ok", "stdout": "class HTTPStream(Stream):\n    def __init__(self, next_chunk, complete):\n        self._next_chunk = next_chunk\n        self._complete = complete\n        self._buffer = io.BytesIO()\n\n    def read(self, size=-1) -> bytes:\n        sections = []\n        length = 0\n\n        # If we have any data in the buffer read that and clear the buffer.\n        buffered = self._buffer.read()\n        if buffered:\n            sections.append(buffered)\n            length += len(buffered)\n            self._buffer.seek(0)\n            self._buffer.truncate(0)\n\n        # Read each chunk in turn.\n        while (size < 0) or (length < size):\n            section = self._next_chunk()\n            sections.append(section)\n            length += len(section)\n            if section == b'':\n                break\n\n        # If we've more data than requested, then push some back into the buffer.\n        output = b''.join(sections)\n        if size > -1 and len(output) > size:\n            output, remainder = output[:size], output[size:]\n            self._buffer.write(remainder)\n            self._buffer.seek(0)\n\n        return output\n\n    def close(self) -> None:\n        self._buffer.close()\n        if self._complete is not None:\n            self._complete()\n\n\nclass MultiPartStream(Stream):\n    def __init__(self, form: list[tuple[str, str]], files: list[tuple[str, str]], boundary=''):\n        self._form = list(form)\n        self._files = list(files)\n        self._boundary = boundary or os.urandom(16).hex()\n        #\u00a0Mutable state...\n        self._form_progress = list(self._form)\n        self._files_progress = list(self._files)\n        self._fin: typing.Any = None\n        self._complete = False\n        self._buffer = io.BytesIO()\n\n    def read(self, size=-1) -> bytes:\n        sections = []\n        length = 0\n\n", "stderr": "", "exit_code": 0}
28run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_streams.py\")\ns = p.read_text()\nold = \"\"\"    def close(self) -> None:\n        self._buffer.close()\n        if self._complete is not None:\n            self._complete()\"\"\"\nnew = \"\"\"    def close(self) -> None:\n        # Server... always read request to completion on keep alives.\n        # We must consume the remaining body before calling _complete.\n        while True:\n            chunk = self.read(64 * 1024)\n            if not chunk:\n                break\n        self._buffer.close()\n        if self._complete is not None:\n            self._complete()\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
29run_command
{
 "command": "sed -n '111,130p' src/httpx/_server.py"
}
{"status": "ok", "stdout": " logging.basicConfig(\n format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n datefmt=\"%Y-%m-%d %H:%M:%S\",\n level=logg
full
{"status": "ok", "stdout": "    logging.basicConfig(\n        format=\"%(levelname)s [%(asctime)s] %(name)s - %(message)s\",\n        datefmt=\"%Y-%m-%d %H:%M:%S\",\n        level=logging.DEBUG\n    )\n\n    backend = NetworkBackend()\n    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n        yield server\n\n\ndef run(app):\n    with serve_http(app) as server:\n        server.wait()\n", "stderr": "", "exit_code": 0}
30run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_server.py\")\ns = p.read_text()\nold = \"\"\"    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n        yield server\"\"\"\nnew = \"\"\"    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        try:\n            server = HTTPServer(server.host, server.port)\n            logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n            yield server\n        except KeyboardInterrupt:\n            pass\n        finally:\n            # Close streams on server exit.\n            # The backend.serve context manager should handle this, but the task says\n            # \"Close streams on server exit\".\n            # Let's see what backend.serve does.\n            pass\"\"\"\n# I need to check NetworkBackend.serve\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
31run_command
{
 "command": "grep -n \"def serve\" src/httpx/_network.py"
}
{"status": "ok", "stdout": "233: def serve(self, host: str, port: int, handler: typing.Callable[[NetworkStream], None]) -> NetworkServer:\n", "stderr": "", "exi
full
{"status": "ok", "stdout": "233:    def serve(self, host: str, port: int, handler: typing.Callable[[NetworkStream], None]) -> NetworkServer:\n", "stderr": "", "exit_code": 0}
32run_command
{
 "command": "sed -n '233,300p' src/httpx/_network.py"
}
{"status": "ok", "stdout": " def serve(self, host: str, port: int, handler: typing.Callable[[NetworkStream], None]) -> NetworkServer:\n listener = self.listen(h
full
{"status": "ok", "stdout": "    def serve(self, host: str, port: int, handler: typing.Callable[[NetworkStream], None]) -> NetworkServer:\n        listener = self.listen(host, port)\n        return NetworkServer(listener, handler)\n\n    def __repr__(self):\n        return \"<NetworkBackend [threaded]>\"\n\n\nSemaphore = threading.Semaphore\nLock = threading.Lock\nsleep = time.sleep\n", "stderr": "", "exit_code": 0}
33run_command
{
 "command": "grep -n \"class NetworkServer\" src/httpx/_network.py"
}
{"status": "ok", "stdout": "156:class NetworkServer:\n", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "156:class NetworkServer:\n", "stderr": "", "exit_code": 0}
34run_command
{
 "command": "sed -n '156,232p' src/httpx/_network.py"
}
{"status": "ok", "stdout": "class NetworkServer:\n def __init__(self, listener: NetworkListener, handler: typing.Callable[[NetworkStream], None]) -> None:\n sel
full
{"status": "ok", "stdout": "class NetworkServer:\n    def __init__(self, listener: NetworkListener, handler: typing.Callable[[NetworkStream], None]) -> None:\n        self.listener = listener\n        self.handler = handler\n        self._max_workers = 5\n        self._executor = None\n        self._thread = None\n        self._streams = list[NetworkStream]\n\n    @property\n    def host(self):\n        return self.listener.host\n\n    @property\n    def port(self):\n        return self.listener.port\n\n    def __enter__(self):\n        self._executor = concurrent.futures.ThreadPoolExecutor(max_workers=self._max_workers)\n        self._executor.submit(self._serve)\n        return self\n\n    def __exit__(self, exc_type, exc_val, exc_tb):\n        self.listener.close()\n        self._executor.shutdown(wait=True)\n\n    def _serve(self):\n        while stream := self.listener.accept():\n            self._executor.submit(self._handler, stream)\n\n    def _handler(self, stream):\n        try:\n            self.handler(stream)\n        finally:\n            stream.close()\n\n\nclass NetworkBackend:\n    def __init__(self, ssl_ctx: ssl.SSLContext | None = None):\n        self._ssl_ctx = self.create_default_context() if ssl_ctx is None else ssl_ctx\n\n    def create_default_context(self) -> ssl.SSLContext:\n        import certifi\n        return ssl.create_default_context(cafile=certifi.where())\n\n    def connect(self, host: str, port: int) -> NetworkStream:\n        \"\"\"\n        Connect to the given address, returning a NetworkStream instance.\n        \"\"\"\n        address = (host, port)\n        timeout = get_current_timeout()\n        sock = socket.create_connection(address, timeout=timeout)\n        return NetworkStream(sock, address)\n\n    def connect_tls(self, host: str, port: int, hostname: str = '') -> NetworkStream:\n        \"\"\"\n        Connect to the given address, returning a NetworkStream instance.\n        \"\"\"\n        address = (host, port)\n  
35run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_server.py\")\ns = p.read_text()\nold = \"\"\"    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        server = HTTPServer(server.host, server.port)\n        logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n        yield server\"\"\"\nnew = \"\"\"    with backend.serve(\"127.0.0.1\", 8080, handler) as server:\n        try:\n            server = HTTPServer(server.host, server.port)\n            logger.info(f\"Serving on {server.url} (Press CTRL+C to quit)\")\n            yield server\n        except KeyboardInterrupt:\n            pass\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
36run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_network.py\")\ns = p.read_text()\nold = \"\"\"    def __exit__(self, exc_type, exc_val, exc_tb):\n        self.listener.close()\n        self._executor.shutdown(wait=True)\"\"\"\nnew = \"\"\"    def __exit__(self, exc_type, exc_val, exc_tb):\n        self.listener.close()\n        # Close streams on server exit.\n        for stream in self._streams:\n            stream.close()\n        self._executor.shutdown(wait=True)\"\"\"\n# Wait, self._streams is initialized as list[NetworkStream] in __init__, but not used.\n# I should update _handler to add streams to self._streams.\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
37run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_network.py\")\ns = p.read_text()\nold = \"\"\"    def _handler(self, stream):\n        try:\n            self.handler(stream)\n        finally:\n            stream.close()\"\"\"\nnew = \"\"\"    def _handler(self, stream):\n        self._streams.append(stream)\n        try:\n            self.handler(stream)\n        finally:\n            stream.close()\n            self._streams.remove(stream)\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
38run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_network.py\")\ns = p.read_text()\nold = \"\"\"    def __exit__(self, exc_type, exc_val, exc_tb):\n        self.listener.close()\n        self._executor.shutdown(wait=True)\"\"\"\nnew = \"\"\"    def __exit__(self, exc_type, exc_val, exc_tb):\n        self.listener.close()\n        for stream in self._streams:\n            stream.close()\n        self._executor.shutdown(wait=True)\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
39run_command
{
 "command": "cd /workspace && python3 - <<'PY'\nfrom pathlib import Path\np = Path(\"src/httpx/_network.py\")\ns = p.read_text()\nold = \"\"\"        self._streams = list[NetworkStream]\"\"\"\nnew = \"\"\"        self._streams: list[NetworkStream] = []\"\"\"\nassert s.count(old) == 1\nt = s.replace(old, new)\np.write_text(t)\nPY\n"
}
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
full
{"status": "ok", "stdout": "", "stderr": "", "exit_code": 0}
40run_command
{
 "command": "git diff"
}
{"status": "ok", "stdout": "diff --git a/src/httpx/_network.py b/src/httpx/_network.py\nindex 5ea9bb5..8554634 100644\n--- a/src/httpx/_network.py\n+++ b/src/ht
full
{"status": "ok", "stdout": "diff --git a/src/httpx/_network.py b/src/httpx/_network.py\nindex 5ea9bb5..8554634 100644\n--- a/src/httpx/_network.py\n+++ b/src/httpx/_network.py\n@@ -160,7 +160,7 @@ class NetworkServer:\n         self._max_workers = 5\n         self._executor = None\n         self._thread = None\n-        self._streams = list[NetworkStream]\n+        self._streams: list[NetworkStream] = []\n \n     @property\n     def host(self):\n@@ -177,6 +177,8 @@ class NetworkServer:\n \n     def __exit__(self, exc_type, exc_val, exc_tb):\n         self.listener.close()\n+        for stream in self._streams:\n+            stream.close()\n         self._executor.shutdown(wait=True)\n \n     def _serve(self):\n@@ -184,10 +186,12 @@ class NetworkServer:\n             self._executor.submit(self._handler, stream)\n \n     def _handler(self, stream):\n+        self._streams.append(stream)\n         try:\n             self.handler(stream)\n         finally:\n             stream.close()\n+            self._streams.remove(stream)\n \n \n class NetworkBackend:\ndiff --git a/src/httpx/_parsers.py b/src/httpx/_parsers.py\nindex 830fccd..c2c37f4 100644\n--- a/src/httpx/_parsers.py\n+++ b/src/httpx/_parsers.py\n@@ -105,6 +105,7 @@ class HTTPParser:\n         # Track connection keep alive...\n         self.send_keep_alive = True\n         self.recv_keep_alive = True\n+        self.keep_alive = True\n \n         # Special states...\n         self.processing_1xx = False\n@@ -375,7 +376,7 @@ class HTTPParser:\n             self.recv_state = State.DONE\n         return body\n \n-    def complete(self):\n+    def reset(self):\n         is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE\n         is_keepalive = self.send_keep_alive and self.recv_keep_alive\n \ndiff --git a/src/httpx/_server.py b/src/httpx/_server.py\nindex 95226d9..4239bb9 100644\n--- a/src/httpx/_server.py\n+++ b/src/httpx/_server.py\n@@ -89,7 +89,7 @@ class HTTPConnection:\n \n     #
41run_command
{
 "command": "submit_patch"
}
{"status": "error", "error_type": "CommandError", "error_message": "/bin/bash: submit_patch: command not found\n", "details": {"stdout": "", "stderr": "/bin/bas
full
{"status": "error", "error_type": "CommandError", "error_message": "/bin/bash: submit_patch: command not found\n", "details": {"stdout": "", "stderr": "/bin/bash: submit_patch: command not found\n", "exit_code": 127}}
42submit_patch
{}
{"status": "ok", "patch_size": 3245, "files_changed": 4}
full
{"status": "ok", "patch_size": 3245, "files_changed": 4}

Patch

diff --git a/src/httpx/_network.py b/src/httpx/_network.py
index 5ea9bb5..8554634 100644
--- a/src/httpx/_network.py
+++ b/src/httpx/_network.py
@@ -160,7 +160,7 @@ class NetworkServer:
         self._max_workers = 5
         self._executor = None
         self._thread = None
-        self._streams = list[NetworkStream]
+        self._streams: list[NetworkStream] = []
 
     @property
     def host(self):
@@ -177,6 +177,8 @@ class NetworkServer:
 
     def __exit__(self, exc_type, exc_val, exc_tb):
         self.listener.close()
+        for stream in self._streams:
+            stream.close()
         self._executor.shutdown(wait=True)
 
     def _serve(self):
@@ -184,10 +186,12 @@ class NetworkServer:
             self._executor.submit(self._handler, stream)
 
     def _handler(self, stream):
+        self._streams.append(stream)
         try:
             self.handler(stream)
         finally:
             stream.close()
+            self._streams.remove(stream)
 
 
 class NetworkBackend:
diff --git a/src/httpx/_parsers.py b/src/httpx/_parsers.py
index 830fccd..c2c37f4 100644
--- a/src/httpx/_parsers.py
+++ b/src/httpx/_parsers.py
@@ -105,6 +105,7 @@ class HTTPParser:
         # Track connection keep alive...
         self.send_keep_alive = True
         self.recv_keep_alive = True
+        self.keep_alive = True
 
         # Special states...
         self.processing_1xx = False
@@ -375,7 +376,7 @@ class HTTPParser:
             self.recv_state = State.DONE
         return body
 
-    def complete(self):
+    def reset(self):
         is_fully_complete = self.send_state == State.DONE and self.recv_state == State.DONE
         is_keepalive = self.send_keep_alive and self.recv_keep_alive
 
diff --git a/src/httpx/_server.py b/src/httpx/_server.py
index 95226d9..4239bb9 100644
--- a/src/httpx/_server.py
+++ b/src/httpx/_server.py
@@ -89,7 +89,7 @@ class HTTPConnection:
 
     # Start it all over again...
     def _complete(self):
-        self._parser.complete
+        self._parser.reset()
         self._idle_expiry = time.monotonic() + self._keepalive_duration
 
 
@@ -116,9 +116,12 @@ def serve_http(endpoint):
 
     backend = NetworkBackend()
     with backend.serve("127.0.0.1", 8080, handler) as server:
-        server = HTTPServer(server.host, server.port)
-        logger.info(f"Serving on {server.url} (Press CTRL+C to quit)")
-        yield server
+        try:
+            server = HTTPServer(server.host, server.port)
+            logger.info(f"Serving on {server.url} (Press CTRL+C to quit)")
+            yield server
+        except KeyboardInterrupt:
+            pass
 
 
 def run(app):
diff --git a/src/httpx/_streams.py b/src/httpx/_streams.py
index c88a63f..9c59ab4 100644
--- a/src/httpx/_streams.py
+++ b/src/httpx/_streams.py
@@ -127,6 +127,12 @@ class HTTPStream(Stream):
         return output
 
     def close(self) -> None:
+        # Server... always read request to completion on keep alives.
+        # We must consume the remaining body before calling _complete.
+        while True:
+            chunk = self.read(64 * 1024)
+            if not chunk:
+                break
         self._buffer.close()
         if self._complete is not None:
             self._complete()

Test output

show
........................                                                 [100%]
24 passed in 0.05s