Skip to content

schedule twisted websocket writes onto the reactor thread - #49

Merged
bentsku merged 2 commits into
mainfrom
twisted-websocket-thread-safety
Sep 30, 2026
Merged

bentsku merged 2 commits into
mainfrom
twisted-websocket-thread-safety

Conversation

@bentsku

@bentsku bentsku commented Sep 30, 2026 •

Copy link
Copy Markdown
Collaborator

Motivation

The twisted websocket listener runs in a threadpool thread, but TwistedWebSocketAdapter called into WebSocketChannel directly, so accept, send, reject and close wrote to the transport (and changed the wsproto state) from that thread. Twisted is not thread-safe, and this breaks connections in two ways:

  • TLS: writes race with the reactor thread processing incoming TLS records on the same OpenSSL connection. The reactor crashes in the TLS layer and drops the connection, or the client receives corrupted frames:
    twisted/protocols/tls.py _flushReceiveBIO -> _flushSendBIO -> bio_read
    OpenSSL.SSL.Error: []
    
    The empty error list comes from OpenSSL's per-thread error queue. Clients see Connection to remote host was lost while reading the upgrade response. This shows up as a flaky wss:// test in LocalStack's CI.
  • Plain TCP: once the messages exceed the socket buffers, writes race with the reactor flushing the same transport buffer, which loses and reorders data. On Linux, streaming 10k × 1 KB messages broke 10 out of 10 connections on main (for example, message 2505 arrived with the contents of message 5240).

Changes

  • TwistedWebSocketAdapter schedules its operations onto the reactor thread:
    • send uses reactor.callFromThread and doesn't wait for the write. Order is kept, since the reactor runs queued calls in order.
    • accept, reject and close use blockingCallFromThread, the same way twisted.web.wsgi writes responses.
    • reject consumes the body in the listener thread, so the reactor doesn't run arbitrary iterators.
  • The adapter gets the reactor from the WSGIResource, and calls directly when it's already on the reactor thread or no reactor is running.
  • Adds twisted[tls] to the dev extra for the new TLS test.

Testing

  • test_websocket_tls_send_and_close_from_listener_thread serves a listener over TLS and runs 200 accept → send → close cycles. It fails every run on main with the error above.
  • test_send_many_messages (twisted and asgi) streams 10k × 1 KB messages over 3 connections and checks each one arrives intact and in order. It fails every run on main for twisted on Linux; asgi passes on both.
  • The rest of the suite is unchanged and green, on macOS and Linux.

Benchmark

Linux (python:3.13-slim, EPollReactor), 32-byte text messages, median of 5 runs. send() is 10k messages streamed from the listener; a cycle is connect → accept → send → close.

ws send() wss send() ws cycle wss cycle
main 166.8k msg/s (loses data under load) crashes (SSL record layer failure) 0.54 ms —
this PR 62.0k msg/s 51.2k msg/s 0.68 ms 1.26 ms
alternative: only schedule on TLS connections 159.3k msg/s (loses data under load) 51.9k msg/s 0.55 ms 1.24 ms
alternative: also wait for every send (blockingCallFromThread) 13.6k msg/s 13.2k msg/s 0.70 ms 1.40 ms

main's plain throughput comes from skipping the thread handoff, which is what loses data. Among the correct variants, not waiting for send is about 4× faster than waiting.

On macOS (SelectReactor), main is much slower: about 300 msg/s over wss://, and the plain ws:// benchmark hung. Writes from another thread don't wake the reactor there. With this PR, both reach about 72k msg/s, with cycles of 0.66 ms (ws) and 1.13 ms (wss).

Benchmark script (python bench_ws.py [--tls])
"""
Benchmark rolo's twisted websocket serving: the throughput of ``send()`` from the listener thread, and the
latency of a full connect/accept/send/close cycle. Run with ``python bench_ws.py [--tls]``.
"""
import argparse
import datetime
import faulthandler
import os
import ssl as stdlib_ssl
import statistics
import threading
import time

import websocket
from cryptography import x509
from cryptography.hazmat.primitives import hashes, serialization
from cryptography.hazmat.primitives.asymmetric import ec
from cryptography.x509.oid import NameOID
from twisted.internet import reactor, ssl
from twisted.protocols.tls import TLSMemoryBIOFactory
from twisted.web.server import Site

from rolo.serving.twisted import (
    HeaderPreservingHTTPChannel,
    HeaderPreservingWSGIResource,
    TwistedRequestAdapter,
    WebsocketResourceDecorator,
)
from rolo.testing.pytest import get_random_tcp_port
from rolo.websocket import WebSocketRequest

MESSAGES = int(os.environ.get("MESSAGES", 10_000))
CYCLES = int(os.environ.get("CYCLES", 1_000))
REPEATS = int(os.environ.get("REPEATS", 5))


def cert_options():
    key = ec.generate_private_key(ec.SECP256R1())
    name = x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, "localhost")])
    now = datetime.datetime.now(datetime.timezone.utc)
    cert = (
        x509.CertificateBuilder()
        .subject_name(name)
        .issuer_name(name)
        .public_key(key.public_key())
        .serial_number(x509.random_serial_number())
        .not_valid_before(now - datetime.timedelta(days=1))
        .not_valid_after(now + datetime.timedelta(days=1))
        .sign(key, hashes.SHA256())
    )
    pem = cert.public_bytes(serialization.Encoding.PEM) + key.private_bytes(
        serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8, serialization.NoEncryption()
    )
    return ssl.PrivateCertificate.loadPEM(pem).options()


@WebSocketRequest.listener
def app(request: WebSocketRequest):
    with request.accept() as ws:
        mode = ws.receive()
        if mode == "stream":
            for _ in range(MESSAGES):
                ws.send("x" * 32)
        else:
            ws.send("hello")


def serve(tls: bool) -> str:
    site = Site(
        WebsocketResourceDecorator(
            original=HeaderPreservingWSGIResource(reactor, reactor.getThreadPool(), None),
            websocketListener=app,
        ),
        requestFactory=TwistedRequestAdapter,
    )
    site.protocol = HeaderPreservingHTTPChannel.protocol_factory
    factory = TLSMemoryBIOFactory(cert_options(), False, site) if tls else site
    port = get_random_tcp_port()
    reactor.listenTCP(port, factory)
    threading.Thread(target=reactor.run, kwargs={"installSignalHandlers": False}, daemon=True).start()
    while not reactor.running:
        time.sleep(0.01)
    return f"{'wss' if tls else 'ws'}://localhost:{port}"


def connect(url: str) -> websocket.WebSocket:
    client = websocket.WebSocket(sslopt={"cert_reqs": stdlib_ssl.CERT_NONE})
    client.connect(url, timeout=10)
    return client


def bench_stream(url: str) -> float:
    client = connect(url)
    start = time.perf_counter()
    client.send("stream")
    for _ in range(MESSAGES):
        client.recv()
    elapsed = time.perf_counter() - start
    client.close()
    return MESSAGES / elapsed


def bench_cycles(url: str) -> float:
    start = time.perf_counter()
    for _ in range(CYCLES):
        client = connect(url)
        client.send("once")
        assert client.recv() == "hello"
        client.close()
    return (time.perf_counter() - start) / CYCLES * 1000


def main():
    faulthandler.dump_traceback_later(60, exit=True)
    parser = argparse.ArgumentParser()
    parser.add_argument("--tls", action="store_true")
    args = parser.parse_args()
    url = serve(args.tls)
    bench_stream(url)  # warm up

    streams = [bench_stream(url) for _ in range(REPEATS)]
    cycles = [bench_cycles(url) for _ in range(REPEATS)]
    print(
        f"{'wss' if args.tls else 'ws '}: send() {statistics.median(streams):>8.0f} msg/s | "
        f"connect/accept/send/close cycle {statistics.median(cycles):.2f} ms (median of {REPEATS})"
    )


if __name__ == "__main__":
    main()

🤖 Generated with Claude Code

The websocket listener runs in a threadpool thread, but wrote to the transport directly. On TLS
connections, this races with the reactor processing incoming TLS records on the same OpenSSL
connection, which breaks the connection (`OpenSSL.SSL.Error: []`) or corrupts frames.

On TLS connections, `send` is now scheduled onto the reactor without waiting for it, and `accept`,
`reject` and `close` wait for their completion. Plain connections keep writing directly.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@bentsku
bentsku marked this pull request as draft September 30, 2026 21:09
…ns too

Writing to the transport from the listener thread also breaks plain connections: it races with the
reactor flushing the same transport, which loses and reorders data once the messages exceed the socket
buffers.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@bentsku bentsku changed the title schedule twisted websocket writes onto the reactor on TLS connections schedule twisted websocket writes onto the reactor thread Sep 30, 2026
@bentsku
bentsku added this pull request to stack #51 September 30, 2026 21:42
@bentsku
bentsku marked this pull request as ready for review September 30, 2026 21:44
@bentsku
bentsku merged commit 428103e into main Sep 30, 2026
5 checks passed
@bentsku
bentsku deleted the twisted-websocket-thread-safety branch September 30, 2026 21:45
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant