Skip to content

FfiClient.request holds the GIL for about 42 ms per MB of request (per-byte ctypes copy) #848

Description

@zdurm

Summary

FfiClient.request copies every request into a new ctypes array one byte at a time:

data = (ctypes.c_ubyte * proto_len)(*proto_data)

data = (ctypes.c_ubyte * proto_len)(*proto_data)

*proto_data unpacks the serialized bytes into one Python int per byte, and the array constructor stores them one at a time. Both loops run in C without releasing the GIL. On my machine this costs about 42 ms per MB, and nothing else in the process runs during it.

Since data streams moved into livekit-ffi (#769), ByteStreamWriter.write puts the whole payload into one byte_stream_write request (data_stream.py#L613-L622). A 4 MB write holds the GIL for about 170 ms; a 16 MB write holds it for about 685 ms. The native call itself takes under 2 ms at 16 MB.

Reproduction

No LiveKit server needed. The script sends a real NewAudioSource request through liblivekit_ffi, padded with an unknown protobuf field that the native decoder skips, so request size is the only variable. A second thread sleeps 1 ms in a loop and records its longest gap.

uv run --no-project --python 3.12 --with livekit==1.1.20 python repro.py
"""FfiClient.request holds the GIL while it copies the request one byte at a time.

No LiveKit server needed. A second thread sleeps 1 ms in a loop and records its
longest gap; a gap far above 1 ms means the main thread held the GIL.

Part 1 times the exact expression from the installed FfiClient.request.
Part 2 sends a real request through liblivekit_ffi (NewAudioSource padded with
an unknown protobuf field, which the native decoder skips), so the size is the
only variable, the same as a large ByteStreamWriter.write.
"""

import ctypes
import inspect
import sys
import threading
import time

from livekit import rtc
from livekit.rtc import _ffi_client
from livekit.rtc._ffi_client import FfiClient, FfiHandle
from livekit.rtc._proto import audio_frame_pb2, ffi_pb2

STOCK_LINE = "data = (ctypes.c_ubyte * proto_len)(*proto_data)"
MB = 1024 * 1024
SIZES_MB = (1, 4, 16)
HEARTBEAT_S = 0.001
STALL_LIMIT_MS = 20.0


class Heartbeat:
    def __init__(self) -> None:
        self.worst = 0.0
        self._stop = threading.Event()
        self._thread = threading.Thread(target=self._run, daemon=True)

    def _run(self) -> None:
        while not self._stop.is_set():
            t = time.perf_counter()
            time.sleep(HEARTBEAT_S)
            self.worst = max(self.worst, time.perf_counter() - t - HEARTBEAT_S)

    def __enter__(self) -> "Heartbeat":
        self._thread.start()
        time.sleep(0.02)
        self.worst = 0.0
        return self

    def __exit__(self, *exc: object) -> None:
        time.sleep(0.02)
        self._stop.set()
        self._thread.join()


def measure(fn) -> tuple[float, float]:
    with Heartbeat() as hb:
        t = time.perf_counter()
        fn()
        elapsed = time.perf_counter() - t
    return elapsed * 1000, hb.worst * 1000


def varint(value: int) -> bytes:
    out = bytearray()
    while True:
        byte, value = value & 0x7F, value >> 7
        out.append(byte | 0x80 if value else byte)
        if not value:
            return bytes(out)


def padded_request(size: int) -> ffi_pb2.FfiRequest:
    req = ffi_pb2.FfiRequest()
    req.new_audio_source.type = audio_frame_pb2.AUDIO_SOURCE_NATIVE
    req.new_audio_source.sample_rate = 48000
    req.new_audio_source.num_channels = 1
    req.new_audio_source.queue_size_ms = 100
    field_1000_bytes = 1000 << 3 | 2
    req.MergeFromString(varint(field_1000_bytes) + varint(size) + b"x" * size)
    return req


def real_request(req: ffi_pb2.FfiRequest) -> None:
    resp = FfiClient.instance.request(req)
    handle = resp.new_audio_source.source.handle.id
    assert handle, "native library rejected the request"
    FfiHandle(handle).dispose()


def main() -> int:
    present = STOCK_LINE in inspect.getsource(_ffi_client)
    print(f"python {sys.version.split()[0]}, livekit-rtc {rtc.__version__}")
    print(f"stock line present in {_ffi_client.__name__}: {present}")

    print("\n1) the stock expression alone, versus a pointer cast over the same bytes")
    for size_mb in SIZES_MB:
        proto_data = b"x" * (size_mb * MB)
        proto_len = len(proto_data)
        stock_ms, stock_stall = measure(lambda: (ctypes.c_ubyte * proto_len)(*proto_data))
        cast_ms, cast_stall = measure(
            lambda: ctypes.cast(ctypes.c_char_p(proto_data), ctypes.POINTER(ctypes.c_ubyte))
        )
        print(
            f"  {size_mb:>2} MB  stock {stock_ms:7.1f} ms, heartbeat stalled {stock_stall:7.1f} ms"
            f"  |  cast {cast_ms:5.2f} ms, heartbeat stalled {cast_stall:5.1f} ms"
        )

    print("\n2) FfiClient.request through liblivekit_ffi (no server)")
    FfiClient.instance  # load the native library outside the timed region
    worst = 0.0
    for size_mb in SIZES_MB:
        req = padded_request(size_mb * MB)
        req_ms, stall = measure(lambda: real_request(req))
        worst = max(worst, stall)
        print(f"  {size_mb:>2} MB  request {req_ms:7.1f} ms, heartbeat stalled {stall:7.1f} ms")

    failed = worst > STALL_LIMIT_MS
    verdict = "FAIL" if failed else "PASS"
    print(
        f"\n{verdict}: FfiClient.request stalled another thread up to {worst:.0f} ms"
        f" (limit {STALL_LIMIT_MS:.0f} ms)"
    )
    return 1 if failed else 0


if __name__ == "__main__":
    sys.exit(main())

Expected

FfiClient.request does not copy the payload in Python, and a large request does not hold the GIL for longer than the native call takes. The native call reads the buffer before it returns, so a pointer to the existing bytes is enough; no copy is needed.

Actual

python 3.12.12, livekit-rtc 1.1.20
stock line present in livekit.rtc._ffi_client: True

1) the stock expression alone, versus a pointer cast over the same bytes
   1 MB  stock    42.4 ms, heartbeat stalled    41.8 ms  |  cast  0.02 ms, heartbeat stalled   0.3 ms
   4 MB  stock   170.3 ms, heartbeat stalled   169.3 ms  |  cast  0.01 ms, heartbeat stalled   0.3 ms
  16 MB  stock   686.4 ms, heartbeat stalled   685.4 ms  |  cast  0.01 ms, heartbeat stalled   0.3 ms

2) FfiClient.request through liblivekit_ffi (no server)
   1 MB  request    43.1 ms, heartbeat stalled    43.1 ms
   4 MB  request   175.8 ms, heartbeat stalled   175.0 ms
  16 MB  request   685.3 ms, heartbeat stalled   684.3 ms

FAIL: FfiClient.request stalled another thread up to 684 ms (limit 20 ms)

The same script with the line replaced by the pointer cast from #847:

python 3.12.12, livekit-rtc 1.1.20
stock line present in livekit.rtc._ffi_client: True

1) the stock expression alone, versus a pointer cast over the same bytes
   1 MB  stock    42.4 ms, heartbeat stalled    42.5 ms  |  cast  0.02 ms, heartbeat stalled   0.3 ms
   4 MB  stock   171.1 ms, heartbeat stalled   170.2 ms  |  cast  0.02 ms, heartbeat stalled   0.3 ms
  16 MB  stock   698.4 ms, heartbeat stalled   697.5 ms  |  cast  0.01 ms, heartbeat stalled   0.3 ms

2) FfiClient.request through liblivekit_ffi (no server)
   1 MB  request     0.2 ms, heartbeat stalled     0.3 ms
   4 MB  request     0.4 ms, heartbeat stalled     0.7 ms
  16 MB  request     1.4 ms, heartbeat stalled     1.6 ms

PASS: FfiClient.request stalled another thread up to 2 ms (limit 20 ms)

Environment

  • livekit 1.1.20 (latest on PyPI), rtc-v1.1.20 at ee527bd
  • The line is unchanged on main at 082e83c (_ffi_client.py line 291). I did not build main from source.
  • CPython 3.12.12, macOS arm64 (Apple M4 Max)
  • Same line in livekit-rtc 1.1.18, where we first hit it.

Impact

Our production worker runs many agent sessions in one process, each on its own event loop. One large byte stream write stalls every session in the process, not only the one that sent it.

Small requests are not affected in practice: audio capture requests carry a pointer to the samples (AudioFrameBufferInfo.data_ptr), not the samples, so they stay small. The stall comes from requests that carry user data, such as byte_stream_write.

Possible fix

Pass a pointer to the serialized bytes instead of copying them:

data = ctypes.cast(ctypes.c_char_p(proto_data), ctypes.POINTER(ctypes.c_ubyte))

proto_data stays referenced until livekit_ffi_request returns, so the pointer stays valid for the call. This is the change in #847, which also removes the per-request c_ubyte_Array_N type that only the cyclic GC can free.

I used an AI assistant to trace this and draft the report; I verified the reproduction myself.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions