Skip to content
Merged
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
56 changes: 56 additions & 0 deletions .github/scripts/sanitizers-macos.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
#!/usr/bin/env bash
#
# Builds the extension with AddressSanitizer and UndefinedBehaviorSanitizer on macOS (Apple Clang, natively), or
# ThreadSanitizer with SANITIZE=thread, and runs the Python tests against it. The Linux counterpart is sanitizers.sh.
#
# .github/scripts/sanitizers-macos.sh [pytest arguments]
#
# The build and a venv without the editable install (whose import hook would load the regular build) are kept
# in build/asan (build/tsan).

set -euo pipefail

SRC="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)"
SANITIZE="${SANITIZE:-address,undefined}"
if [ "$SANITIZE" = thread ]; then
BUILD="${WRTC_SANITIZERS_BUILD_DIR:-$SRC/build/tsan}"
else
BUILD="${WRTC_SANITIZERS_BUILD_DIR:-$SRC/build/asan}"
fi
if [ -z "${PYTHON:-}" ]; then
uv python install -q 3.13
PYTHON="$(uv python find 3.13)"
fi

if [ ! -x "$BUILD/venv/bin/python" ]; then
uv venv -q "$BUILD/venv" --python "$PYTHON"
uv pip install -q --python "$BUILD/venv/bin/python" cmake ninja "pybind11>=3.0" pytest pytest-asyncio pytest-timeout
fi
export PATH="$BUILD/venv/bin:$PATH"

cmake -S "$SRC" -B "$BUILD" -G Ninja \
-DCMAKE_BUILD_TYPE=RelWithDebInfo \
-DWRTC_SANITIZE="$SANITIZE" \
-DPython_EXECUTABLE="$BUILD/venv/bin/python" \
-Dpybind11_DIR="$("$BUILD/venv/bin/python" -m pybind11 --cmakedir)" > /dev/null
cmake --build "$BUILD"

# The interpreter isn't instrumented: the runtime must be loaded before anything else
if [ "$SANITIZE" = thread ]; then
DYLD_INSERT_LIBRARIES="$(clang -print-runtime-dir)/libclang_rt.tsan_osx_dynamic.dylib"
# libwebrtc isn't instrumented: races of its own aren't reported
export TSAN_OPTIONS="halt_on_error=1:report_signal_unsafe=0:strip_env=0"
else
DYLD_INSERT_LIBRARIES="$(clang -print-runtime-dir)/libclang_rt.asan_osx_dynamic.dylib"
export PYTHONMALLOC=malloc
# LeakSanitizer isn't supported on macOS; strip_env=0 keeps the runtime in subprocesses of the tests;
# container annotations can't match libwebrtc's libc++, whose code the linker may fold with ours
export ASAN_OPTIONS="detect_leaks=0:detect_container_overflow=0:halt_on_error=1:abort_on_error=0:strict_init_order=1:strip_env=0"
export UBSAN_OPTIONS="print_stacktrace=1:halt_on_error=1"
fi
export DYLD_INSERT_LIBRARIES
export PYTHONPATH="$BUILD/python-webrtc/cpp:$SRC/python-webrtc/python:$SRC"
export PYTHONDONTWRITEBYTECODE=1

cd "$SRC"
"$BUILD/venv/bin/python" -m pytest tests --ignore=tests/wpt -p no:cacheprovider --capture=sys "$@"
18 changes: 16 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -84,14 +84,28 @@ jobs:
name: sanitizers (ASan + UBSan)
runs-on: ubuntu-latest
container: quay.io/pypa/manylinux_2_28_x86_64
timeout-minutes: 30
timeout-minutes: 60
steps:
- uses: actions/checkout@v7
- run: .github/scripts/sanitizers.sh
# the collector running on libwebrtc threads, whenever they emit (see tests/conftest.py)
- run: .github/scripts/sanitizers.sh --gc-on-emit
# long random sequences of calls (tests/chaos.py)
- run: .github/scripts/sanitizers.sh --stress -k "stress or chaos"

sanitizers-macos:
name: sanitizers (ASan + UBSan, macOS)
runs-on: macos-15
timeout-minutes: 60
steps:
- uses: actions/checkout@v7
- uses: astral-sh/setup-uv@v7
- run: .github/scripts/sanitizers-macos.sh
- run: .github/scripts/sanitizers-macos.sh --gc-on-emit

publish:
if: startsWith(github.ref, 'refs/tags/v')
needs: [lint, sdist, wheels, sanitizers]
needs: [lint, sdist, wheels, sanitizers, sanitizers-macos]
runs-on: ubuntu-latest
environment: pypi
permissions:
Expand Down
9 changes: 8 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
.PHONY: dev test lint format stub wheels doc clean
.PHONY: dev test asan tsan lint format stub wheels doc clean

# editable install; the extension is rebuilt automatically on import after C++ changes
dev:
Expand All @@ -9,6 +9,13 @@ dev:
test:
uv run --no-sync pytest tests $(O)

# the tests against an ASan+UBSan build, natively on macOS (on Linux: .github/scripts/sanitizers.sh)
asan:
.github/scripts/sanitizers-macos.sh $(O)

tsan:
SANITIZE=thread .github/scripts/sanitizers-macos.sh $(O)

lint:
uvx ruff check
uvx ruff format --check
Expand Down
30 changes: 30 additions & 0 deletions python-webrtc/cpp/src/interfaces/interfaces.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,16 @@
#include "rtc_data_channel.h"
#include "rtc_dtmf_sender.h"
#include "rtc_peer_connection.h"
#include "../media/media_stream_track_processor.h"
#include "../media/track_generator.h"
#include "../media/video_frame_buffer.h"
#include "../utils/alive_count.h"
#include "../utils/gil.h"

#include <map>
#include <string>

#include <pybind11/stl.h>

namespace python_webrtc {

Expand All @@ -35,5 +45,25 @@ namespace python_webrtc {
RTCRtpTransceiver::Init(m);
RTCDataChannel::Init(m);
RTCPeerConnection::Init(m);

// the native objects alive, by type, so tests can check that none leaks
m.def("_alive", []() {
return std::map<std::string, int>{
{"RTCPeerConnection", AliveCount<RTCPeerConnection>::count.load()},
{"MediaStreamTrack", MediaStreamTrack::holder().Alive()},
{"MediaStream", MediaStream::holder().Alive()},
{"RTCRtpTransceiver", RTCRtpTransceiver::holder().Alive()},
{"RTCRtpSender", RTCRtpSender::holder().Alive()},
{"RTCRtpReceiver", RTCRtpReceiver::holder().Alive()},
{"RTCDTMFSender", RTCDTMFSender::holder().Alive()},
{"RTCDataChannel", RTCDataChannel::holder().Alive()},
{"RTCSctpTransport", RTCSctpTransport::holder().Alive()},
{"RTCDtlsTransport", RTCDtlsTransport::holder().Alive()},
{"RTCIceTransport", RTCIceTransport::holder().Alive()},
{"MediaStreamTrackProcessor", AliveCount<MediaStreamTrackProcessor>::count.load()},
{"TrackGenerator", AliveCount<TrackGenerator>::count.load()},
{"VideoFrameBuffer", AliveCount<VideoFrameBuffer>::count.load()},
};
}, nogil());
}
}
44 changes: 4 additions & 40 deletions python-webrtc/cpp/src/interfaces/media_stream.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -68,37 +68,14 @@ namespace python_webrtc {
}

std::vector<std::shared_ptr<MediaStreamTrack>> MediaStream::SyncTracks() {
auto tracks = std::vector<std::shared_ptr<MediaStreamTrack>>();
decltype(_tracks) current;
// read before locking: they're calls to the signaling thread, where OnChanged takes the lock
auto streamTracks = this->tracks();
{
std::lock_guard<std::mutex> lock(_tracksMutex);
for (auto const &track: streamTracks) {
auto it = _tracks.find(track.get());
auto wrapper = it != _tracks.end() ? it->second : MediaStreamTrack::holder().GetOrCreate(_factory, track);
current[track.get()] = wrapper;
tracks.push_back(std::move(wrapper));
}
std::swap(_tracks, current);
// Python keeps the wrappers, the holder finds them
std::vector<std::shared_ptr<MediaStreamTrack>> tracks;
for (auto const &track: this->tracks()) {
tracks.push_back(MediaStreamTrack::holder().GetOrCreate(_factory, track));
}
// wrappers of removed tracks are released here, out of the lock
return tracks;
}

std::shared_ptr<MediaStreamTrack> MediaStream::WrapTrack(
webrtc::scoped_refptr<webrtc::MediaStreamTrackInterface> track) {
std::lock_guard<std::mutex> lock(_tracksMutex);
auto it = _tracks.find(track.get());
if (it != _tracks.end()) {
return it->second;
}

auto wrapper = MediaStreamTrack::holder().GetOrCreate(_factory, track);
_tracks[track.get()] = wrapper;
return wrapper;
}

webrtc::scoped_refptr<webrtc::MediaStreamInterface> MediaStream::stream() {
return _stream;
}
Expand Down Expand Up @@ -182,9 +159,6 @@ namespace python_webrtc {
} else {
_stream->AddTrack(static_cast<webrtc::scoped_refptr<webrtc::VideoTrackInterface>>(*mediaStreamTrack));
}

std::lock_guard<std::mutex> lock(_tracksMutex);
_tracks[track.get()] = mediaStreamTrack;
}

void MediaStream::RemoveTrack(MediaStreamTrack &mediaStreamTrack) {
Expand All @@ -200,16 +174,6 @@ namespace python_webrtc {
} else {
_stream->RemoveTrack(static_cast<webrtc::scoped_refptr<webrtc::VideoTrackInterface>>(mediaStreamTrack));
}

std::shared_ptr<MediaStreamTrack> removed;
{
std::lock_guard<std::mutex> lock(_tracksMutex);
auto it = _tracks.find(track.get());
if (it != _tracks.end()) {
removed = std::move(it->second);
_tracks.erase(it);
}
}
}

std::shared_ptr<MediaStream> MediaStream::Clone() {
Expand Down
2 changes: 0 additions & 2 deletions python-webrtc/cpp/src/interfaces/media_stream.h
Original file line number Diff line number Diff line change
Expand Up @@ -78,13 +78,11 @@ namespace python_webrtc {
// wrappers of the current tracks of the stream; the stream owns them, so their state outlives Python references
std::vector<std::shared_ptr<MediaStreamTrack>> SyncTracks();

std::shared_ptr<MediaStreamTrack> WrapTrack(webrtc::scoped_refptr<webrtc::MediaStreamTrackInterface>);

std::shared_ptr<PeerConnectionFactory> _factory;
webrtc::scoped_refptr<webrtc::MediaStreamInterface> _stream;

std::mutex _tracksMutex;
std::unordered_map<webrtc::MediaStreamTrackInterface *, std::shared_ptr<MediaStreamTrack>> _tracks;
// the tracks the stream had when last notified, or as Python changed them, guarded by _tracksMutex
std::set<webrtc::MediaStreamTrackInterface *> _known;

Expand Down
15 changes: 9 additions & 6 deletions python-webrtc/cpp/src/interfaces/media_stream_track.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -42,11 +42,14 @@ namespace python_webrtc {

_track = nullptr;
// released as the listeners are (see DropListeners)
if (!PythonAlive()) {
(void) _constraints.release();
} else if (_constraints) {
pybind11::gil_scoped_acquire gil;
pybind11::object dropped = std::move(_constraints);
{
PythonEntry entry;
if (!entry) {
(void) _constraints.release();
} else if (_constraints) {
pybind11::gil_scoped_acquire gil;
pybind11::object dropped = std::move(_constraints);
}
}
DropListeners();
}
Expand Down Expand Up @@ -257,7 +260,7 @@ namespace python_webrtc {
std::optional<std::tuple<int, int, double>> camera;
bool microphone = false;
{
pybind11::gil_scoped_release release;
gil_release release;
video = _monitor.video();
audio = _monitor.audio();
camera = GetCamera();
Expand Down
46 changes: 44 additions & 2 deletions python-webrtc/cpp/src/interfaces/peer_connection_factory.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,9 @@

#include "peer_connection_factory.h"
#include "../media/playout_audio_device.h"
#include "../media/wakeup.h"
#include "../utils/gil.h"
#include "../utils/instance_holder.h"
#include "../utils/libwebrtc_thread.h"

#include <api/create_peerconnection_factory.h>
Expand All @@ -24,8 +26,13 @@
#include <api/video_codecs/video_encoder_factory_template_libvpx_vp9_adapter.h>
#include <rtc_base/ssl_adapter.h>

#include <stdexcept>
#include <thread>

#ifndef _WIN32
#include <pthread.h>
#endif

namespace python_webrtc {

// Royalty-free codecs only (the prebuilts have no H.264).
Expand All @@ -43,7 +50,7 @@ namespace python_webrtc {
std::mutex PeerConnectionFactory::_mutex{};
std::atomic<int> PeerConnectionFactory::_alive{0};

PeerConnectionFactory::PeerConnectionFactory() {
PeerConnectionFactory::PeerConnectionFactory() : _generation(forks.load()) {
_alive++;

_workerThread = webrtc::Thread::CreateWithSocketServer();
Expand Down Expand Up @@ -108,7 +115,19 @@ namespace python_webrtc {
_alive--;
}

void RunOnSignalingThread(PeerConnectionFactory &factory, const std::function<void()> &function) {
gil_release_if_held release;
factory._signalingThread->BlockingCall([&]() { function(); });
}

std::shared_ptr<PeerConnectionFactory> PeerConnectionFactory::Create() {
#ifdef __APPLE__
// libwebrtc runs its task queues on libdispatch, which crashes in the child of a fork
if (forks.load() > 0) {
throw std::runtime_error("python-webrtc can't be used in the child of a fork on macOS (libdispatch doesn't support "
"it): use the spawn start method of multiprocessing");
}
#endif
return {new PeerConnectionFactory(), &PeerConnectionFactory::Destroy};
}

Expand All @@ -123,6 +142,10 @@ namespace python_webrtc {
}

void PeerConnectionFactory::Destroy(PeerConnectionFactory *factory) {
// leaked while the interpreter finalizes or in a forked child: its threads may hang or be gone
if (!PythonAlive() || factory->_generation != forks.load()) {
return;
}
// the last owner may be released by a task on one of the factory threads, which can't stop itself
if (factory->_workerThread->IsCurrent() || factory->_signalingThread->IsCurrent()) {
std::thread([factory]() { delete factory; }).detach();
Expand All @@ -140,8 +163,27 @@ namespace python_webrtc {
[[maybe_unused]] bool result = webrtc::InitializeSSL();
assert(result);

#ifndef _WIN32
// libwebrtc threads don't survive a fork: the child forgets the factories (also runs before exec, keep it minimal)
pthread_atfork(
[]() {
_mutex.lock();
Wakeup::LockForFork();
},
[]() {
Wakeup::UnlockAfterFork();
_mutex.unlock();
},
[]() {
forks++;
_default.reset();
Wakeup::UnlockAfterFork();
_mutex.unlock();
});
#endif

pybind11::class_<PeerConnectionFactory, std::shared_ptr<PeerConnectionFactory>>(m, "PeerConnectionFactory")
.def(pybind11::init(&PeerConnectionFactory::Create), nogil())
.def(pybind11::init(nogil_factory(&PeerConnectionFactory::Create)))
.def_static("getOrCreateDefault", &PeerConnectionFactory::GetOrCreateDefault, nogil())
.def_static("dispose", &PeerConnectionFactory::Dispose, nogil());

Expand Down
3 changes: 3 additions & 0 deletions python-webrtc/cpp/src/interfaces/peer_connection_factory.h
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,9 @@ namespace python_webrtc {
std::unique_ptr<webrtc::Thread> _workerThread;

private:
// of the process the factory was created in (see forks)
const int _generation;

static void Destroy(PeerConnectionFactory *);

static std::weak_ptr<PeerConnectionFactory> _default;
Expand Down
10 changes: 5 additions & 5 deletions python-webrtc/cpp/src/interfaces/rtc_data_channel.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ namespace python_webrtc {
}
// messages received meanwhile are delivered to the observer once it's registered, and are held too
_channel->RegisterObserver(this);
holder().SetObserver(_channel.get(), this);
// a channel announced by the remote peer is open already, but its open event follows the datachannel one
if (_lastState == DataState::kOpen) {
Emit("open", _lastState);
Expand All @@ -39,11 +40,10 @@ namespace python_webrtc {
RTCDataChannel::~RTCDataChannel() {
BlockingDestructor release("RTCDataChannel");

// the channel has a single observer slot, a newer wrapper of it may have taken it over already
auto replaced = holder().HasLive(_channel.get());
// callbacks run on the signaling thread, so after this none of them can be running or start again
_factory->_signalingThread->BlockingCall([this, replaced]() {
if (!replaced) {
_factory->_signalingThread->BlockingCall([this]() {
// a newer wrapper of the channel may have taken its single observer slot
if (holder().TakeObserver(_channel.get(), this)) {
_channel->UnregisterObserver();
}
});
Expand Down Expand Up @@ -80,7 +80,7 @@ namespace python_webrtc {
.def("close", &RTCDataChannel::Close, nogil())
.def("_surfaceState", &RTCDataChannel::SurfaceState, nogil(), pybind11::arg("state"))
.def("_decreaseBufferedAmount", &RTCDataChannel::DecreaseBufferedAmount, nogil(), pybind11::arg("sent"))
.def("_release", &RTCDataChannel::Release);
.def("_release", &RTCDataChannel::Release, nogil());
}

InstanceHolder<RTCDataChannel, webrtc::DataChannelInterface> &RTCDataChannel::holder() {
Expand Down
Loading
Loading