Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions packages/runtime-sdk/README.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,14 @@
# workers-runtime-sdk

Runtime SDK for Python Cloudflare Workers.

## ASGI response encoding

The ASGI adapter preserves application response bytes and `Content-Encoding`.
Applications and middleware are responsible for encoding their response bodies;
the adapter uses the Workers Fetch API's manual encoding mode to avoid
compressing an already-compressed response a second time.

Applications must negotiate supported encodings with the request's
`Accept-Encoding` header. Manual encoding does not make a gzip response
appropriate for an identity-only client.
11 changes: 9 additions & 2 deletions packages/runtime-sdk/src/workers/asgi.py
Original file line number Diff line number Diff line change
Expand Up @@ -261,8 +261,12 @@ async def send(got):
transform_stream = TransformStream.new()
readable = transform_stream.readable
writer = transform_stream.writable.getWriter()
# ASGI bytes already match Content-Encoding; do not encode twice.
resp = Response.new(
readable, headers=_to_js_headers(headers), status=status
readable,
headers=_to_js_headers(headers),
status=status,
encodeBody="manual",
)
result.set_result(resp)
with acquire_js_buffer(body) as jsbytes:
Expand All @@ -281,7 +285,10 @@ async def send(got):
buf = px.getBuffer()
px.destroy()
resp = Response.new(
buf.data, headers=_to_js_headers(headers), status=status
buf.data,
headers=_to_js_headers(headers),
status=status,
encodeBody="manual",
)
result.set_result(resp)
finished_response.set()
Expand Down
15 changes: 15 additions & 0 deletions packages/runtime-sdk/tests/workerd-test/asgi/asgi.wd-test
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,25 @@ const unitTests :Workerd.Config = (
],
bindings = [
( name = "SELF", service = "python-asgi" ),
( name = "ASGI_HTTP", service = "asgi-http" ),
],
compatibilityDate = "%COMPAT_DATE",
compatibilityFlags = ["python_workers"],
)
),
( name = "asgi-http",
external = (
address = "loopback:asgi-http",
http = (),
)
)
],
sockets = [
(
name = "asgi-http",
address = "loopback:asgi-http",
http = (),
service = "python-asgi",
),
],
);
38 changes: 37 additions & 1 deletion packages/runtime-sdk/tests/workerd-test/asgi/tests/test_asgi.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,12 @@
import js
import pytest
from pyodide.ffi import to_js
from worker import STREAMING_CHUNK_SIZE, STREAMING_NUM_CHUNKS, example_hdr
from worker import (
ENCODED_RESPONSE_PAYLOAD,
STREAMING_CHUNK_SIZE,
STREAMING_NUM_CHUNKS,
example_hdr,
)

import asgi
from workers import Request, env
Expand Down Expand Up @@ -89,6 +94,37 @@ async def test_streaming():
assert all(b == i % 256 for b in chunk)


@pytest.mark.asyncio
@pytest.mark.parametrize("encoding", ["gzip", "identity"])
@pytest.mark.parametrize("mode", ["buffered", "streaming"])
async def test_already_encoded_response_crosses_http_unchanged(encoding, mode):
path = f"/encoded/{encoding}/{mode}"
# SELF.fetch does not serialize the response over HTTP, so it cannot catch
# workerd compressing already-gzipped ASGI bytes a second time. ASGI_HTTP
# instead connects to the worker's loopback HTTP socket.
response = await env.ASGI_HTTP.fetch(
f"http://example.com{path}",
headers={"accept-encoding": "gzip, identity"},
)

assert response.status == 201
assert response.headers["content-type"] == "application/json; charset=utf-8"
assert response.headers["content-encoding"] == encoding
assert response.headers["vary"] == "Accept-Encoding"
assert response.headers["x-asgi-response"] == "already-encoded"

reader = response.body.getReader()
chunks = []
while True:
result = await reader.read()
if result.done:
break
chunks.append(result.value.to_bytes())
# Fetch decodes the HTTP content encoding once, just as a browser does.
# Double gzip therefore leaves a gzip member here instead of the JSON.
assert b"".join(chunks) == ENCODED_RESPONSE_PAYLOAD


class _ListHandler(logging.Handler):
"""A logging handler that captures records into a list for assertions."""

Expand Down
69 changes: 69 additions & 0 deletions packages/runtime-sdk/tests/workerd-test/asgi/worker.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import gzip
import os

from testlib.entrypoint import run_pytest
Expand Down Expand Up @@ -159,6 +160,63 @@ async def __call__(self, scope, receive, send):
)


ENCODED_RESPONSE_PAYLOAD = (
b'{"message":"ASGI middleware already encoded this response",'
b'"padding":"' + b"compression-regression-" * 128 + b'"}'
)


class EncodedResponseApp:
"""Return middleware-produced bytes unchanged, buffered or streamed."""

def __init__(self, encoding, streaming):
self.encoding = encoding
self.streaming = streaming
self.body = (
gzip.compress(ENCODED_RESPONSE_PAYLOAD, mtime=0)
if encoding == "gzip"
else ENCODED_RESPONSE_PAYLOAD
)

async def __call__(self, scope, receive, send):
if scope["type"] == "lifespan":
message = await receive()
if message["type"] == "lifespan.startup":
await send({"type": "lifespan.startup.complete"})
return

await receive()
headers = [
(b"content-type", b"application/json; charset=utf-8"),
(b"content-encoding", self.encoding.encode()),
(b"vary", b"Accept-Encoding"),
(b"x-asgi-response", b"already-encoded"),
]
if not self.streaming:
headers.append((b"content-length", str(len(self.body)).encode()))
await send(
{
"type": "http.response.start",
"status": 201,
"headers": headers,
}
)
if self.streaming:
# Split inside the gzip header and before its footer, rather than
# emitting separate gzip members for the individual ASGI chunks.
chunks = (self.body[:7], self.body[7:-8], self.body[-8:], b"")
for index, chunk in enumerate(chunks):
await send(
{
"type": "http.response.body",
"body": chunk,
"more_body": index < len(chunks) - 1,
}
)
else:
await send({"type": "http.response.body", "body": self.body})


# ---------------------------------------------------------------------------
# App instances and constants
# ---------------------------------------------------------------------------
Expand Down Expand Up @@ -244,6 +302,13 @@ async def __call__(self, scope, receive, send):
scope_echo_app = ScopeEchoApp()
late_failure_stream_app = LateFailureStreamApp()
multi_cookie_app = MultiCookieApp()
encoded_response_apps = {
f"/encoded/{encoding}/{mode}": EncodedResponseApp(
encoding, streaming=mode == "streaming"
)
for encoding in ("gzip", "identity")
for mode in ("buffered", "streaming")
}

example_hdr = {"Header1": "Value1", "Header2": "Value2"}

Expand All @@ -255,6 +320,10 @@ async def fetch(self, request):
url = URL.new(request.url)
path = url.pathname

if path in encoded_response_apps:
return await asgi.fetch(
encoded_response_apps[path], request, self.env, self.ctx
)
if path == "/sse":
return await asgi.fetch(sse_app, request, self.env, self.ctx)
elif path == "/stream":
Expand Down
Loading