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.
Summary
FfiClient.requestcopies every request into a new ctypes array one byte at a time:python-sdks/livekit-rtc/livekit/rtc/_ffi_client.py
Line 291 in ee527bd
*proto_dataunpacks 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.writeputs the whole payload into onebyte_stream_writerequest (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
NewAudioSourcerequest throughliblivekit_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.Expected
FfiClient.requestdoes 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 existingbytesis enough; no copy is needed.Actual
The same script with the line replaced by the pointer cast from #847:
Environment
rtc-v1.1.20at ee527bdmainat 082e83c (_ffi_client.pyline 291). I did not buildmainfrom source.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 asbyte_stream_write.Possible fix
Pass a pointer to the serialized bytes instead of copying them:
proto_datastays referenced untillivekit_ffi_requestreturns, so the pointer stays valid for the call. This is the change in #847, which also removes the per-requestc_ubyte_Array_Ntype that only the cyclic GC can free.I used an AI assistant to trace this and draft the report; I verified the reproduction myself.