From 3f18e7aaea129fb5dce5c1c0d1cd6c3716a969ce Mon Sep 17 00:00:00 2001 From: tintinhamans <5984296+tintinhamans@users.noreply.github.com> Date: Fri, 25 Sep 2026 05:03:26 +0200 Subject: [PATCH] steamwebrtc: Let resolver threads exit after the lookup --- src/external/steamwebrtc/ice_session.cpp | 186 ++++++++++++++++++++++- 1 file changed, 185 insertions(+), 1 deletion(-) diff --git a/src/external/steamwebrtc/ice_session.cpp b/src/external/steamwebrtc/ice_session.cpp index 57d2e2b5..591d039f 100644 --- a/src/external/steamwebrtc/ice_session.cpp +++ b/src/external/steamwebrtc/ice_session.cpp @@ -5,16 +5,22 @@ #include #include +#include +#include #include +#include #include #include #include #include #include +#include #include +#include #include #include +#include #include #include @@ -55,6 +61,182 @@ extern "C" STEAMWEBRTC_DECLSPEC IICESession *CreateWebRTCICESession( const ICESessionConfig &cfg, IICESessionDelegate *pDelegate, int nInterfaceVersion ); } +namespace rtc +{ + // Defined in net_helpers.cc, but not declared in a header. + int ResolveHostname( const std::string &hostname, int family, std::vector *addresses ); +} + +namespace { + +struct AsyncResolverState +{ + std::mutex mutex; + bool alive = true; // Guarded by mutex + std::atomic finished{ false }; +}; + +using ResolverThreads = std::vector< std::pair< std::thread, std::shared_ptr > >; + +// Leaked on purpose: joinable threads would terminate at exit. +std::mutex s_mutexResolverThreads; +ResolverThreads &s_vecResolverThreads = *new ResolverThreads; + +void AbandonResolverThread( std::thread thread, std::shared_ptr state ) +{ + std::lock_guard lock( s_mutexResolverThreads ); + for ( auto it = s_vecResolverThreads.begin(); it != s_vecResolverThreads.end(); ) + { + if ( it->second->finished ) + { + it->first.join(); + it = s_vecResolverThreads.erase( it ); + } + else + { + ++it; + } + } + s_vecResolverThreads.emplace_back( std::move( thread ), std::move( state ) ); +} + +void JoinResolverThreads() +{ + ResolverThreads threads; + { + std::lock_guard lock( s_mutexResolverThreads ); + threads.swap( s_vecResolverThreads ); + } + for ( auto &thread : threads ) + thread.first.join(); +} + +//----------------------------------------------------------------------------- +// rtc::AsyncResolver's thread idles until the resolver is destroyed (webrtc:12659) +//----------------------------------------------------------------------------- +class CAsyncResolver final : public rtc::AsyncResolverInterface +{ +public: + void Start( const rtc::SocketAddress &addr ) override + { + Assert( !worker_.joinable() ); + addr_ = addr; + + thread_ = rtc::Thread::Current(); + Assert( thread_ != nullptr ); + + CAsyncResolver *pResolver = this; + std::shared_ptr state = state_; + rtc::Thread *pCallingThread = thread_; + std::string hostname = addr.hostname(); + int family = addr.family(); + + worker_ = std::thread( [pResolver, state, pCallingThread, hostname, family]() + { + std::vector addresses; + int error = rtc::ResolveHostname( hostname, family, &addresses ); + + // The calling thread may be gone once Destroy() has run. + std::lock_guard lock( state->mutex ); + state->finished = true; + if ( !state->alive ) + return; + pCallingThread->PostTask( webrtc::ToQueuedTask( + [pResolver, state, addresses = std::move( addresses ), error]() mutable + { + { + std::lock_guard lock( state->mutex ); + if ( !state->alive ) + return; + } + pResolver->ResolveDone( std::move( addresses ), error ); + } ) ); + } ); + } + + bool GetResolvedAddress( int family, rtc::SocketAddress *addr ) const override + { + if ( error_ != 0 || addresses_.empty() ) + return false; + + *addr = addr_; + for ( const rtc::IPAddress &ip : addresses_ ) + { + if ( family == ip.family() ) + { + addr->SetResolvedIP( ip ); + return true; + } + } + return false; + } + + int GetError() const override { return error_; } + + void Destroy( bool wait ) override + { + // The posted result runs on this thread too, so it can't race with Destroy(). + Assert( thread_ == nullptr || thread_->IsCurrent() ); + + { + std::lock_guard lock( state_->mutex ); + state_->alive = false; + } + + if ( worker_.joinable() ) + { + if ( wait ) + worker_.join(); + else + AbandonResolverThread( std::move( worker_ ), state_ ); + } + + // Called from a SignalDone handler: ResolveDone() deletes us after the signal. + if ( in_signal_ ) + { + destroy_pending_ = true; + return; + } + delete this; + } + +private: + ~CAsyncResolver() override = default; + + void ResolveDone( std::vector addresses, int error ) + { + addresses_ = std::move( addresses ); + error_ = error; + in_signal_ = true; + SignalDone( this ); + in_signal_ = false; + if ( destroy_pending_ ) + delete this; + } + + rtc::SocketAddress addr_; + rtc::Thread *thread_ = nullptr; + std::vector addresses_; + int error_ = -1; + bool in_signal_ = false; + bool destroy_pending_ = false; + std::thread worker_; + std::shared_ptr state_ = std::make_shared(); +}; + +class CPacketSocketFactory final : public rtc::BasicPacketSocketFactory +{ +public: + explicit CPacketSocketFactory( rtc::Thread *thread ) : rtc::BasicPacketSocketFactory( thread ) {} + + rtc::AsyncResolverInterface *CreateAsyncResolver() override + { + return new CAsyncResolver(); + } +}; + +} // namespace + //----------------------------------------------------------------------------- // Class to represent an ICE connection //----------------------------------------------------------------------------- @@ -201,6 +383,8 @@ CICESession::~CICESession() delete s_pSocketServer; s_pSocketServer = nullptr; + JoinResolverThreads(); + rtc::CleanupSSL(); } s_mutex.unlock(); @@ -240,7 +424,7 @@ bool CICESession::BInitializeOnSocketThread( const ICESessionConfig &cfg ) default_network_manager_.reset(new rtc::BasicNetworkManager()); default_socket_factory_.reset( - new rtc::BasicPacketSocketFactory( s_pSocketThread )); + new CPacketSocketFactory( s_pSocketThread )); webrtc::TurnCustomizer *turn_customizer = nullptr;