From 8d20e35f02ebacdaf5f76d29f8dcadb9c0af4165 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Sun, 8 May 2022 22:14:16 +0200 Subject: [PATCH] Non blocking SignalClient --- .gitignore | 7 +- .gitmodules | 2 +- CMakeLists.txt | 9 ++- main.cpp | 14 ++-- protocol | 2 +- src/proto/livekit_models.pb.cc | 117 ++++++++++++++++++++-------- src/proto/livekit_models.pb.h | 62 +++++++++++++++ src/proto/livekit_rtc.pb.cc | 107 +++++++++++++++++++------- src/proto/livekit_rtc.pb.h | 67 ++++++++++++++++ src/room.cpp | 5 ++ src/room.h | 27 +++++++ src/rtc_engine.cpp | 14 ++++ src/rtc_engine.h | 25 ++++++ src/signal_client.cpp | 135 ++++++++++++++++++++++++++++----- src/signal_client.h | 39 +++++++++- src/utils.h | 36 +++++++++ 16 files changed, 575 insertions(+), 93 deletions(-) create mode 100644 src/room.cpp create mode 100644 src/room.h create mode 100644 src/rtc_engine.cpp create mode 100644 src/rtc_engine.h create mode 100644 src/utils.h diff --git a/.gitignore b/.gitignore index a5a842b..0f3263f 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,7 @@ cmake-build-debug -.idea/* \ No newline at end of file +.idea +build +.vscode + +# atm +thirdparty \ No newline at end of file diff --git a/.gitmodules b/.gitmodules index 0ff103d..1752bd7 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,3 +1,3 @@ [submodule "protocol"] path = protocol - url = https://github.com/livekit/protocol + url = https://github.com/livekit/protocol \ No newline at end of file diff --git a/CMakeLists.txt b/CMakeLists.txt index 37ce467..e49b4bc 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -7,8 +7,13 @@ find_package(Boost COMPONENTS system REQUIRED) include_directories(${Boost_INCLUDE_DIRS}) find_package(Protobuf REQUIRED) -include_directories(${Protobuf_INCLUDE_DIR}) +include_directories(${Protobuf_INCLUDE_DIRS}) + +find_package(spdlog CONFIG REQUIRED) + +include_directories(thirdparty/webrtc/include) + file(GLOB_RECURSE SRC src/*) add_executable(livekit-native ${SRC} main.cpp) -target_link_libraries(livekit-native ${Boost_LIBRARIES} ${Protobuf_LIBRARIES}) \ No newline at end of file +target_link_libraries(livekit-native PRIVATE ${Boost_LIBRARIES} ${Protobuf_LIBRARIES} spdlog::spdlog) diff --git a/main.cpp b/main.cpp index 543083f..8b418c1 100644 --- a/main.cpp +++ b/main.cpp @@ -1,12 +1,16 @@ #include -#include #include "src/signal_client.h" +#include "spdlog/spdlog.h" int main() { - std::cout << "Hello, World!" << std::endl; + spdlog::info("Starting LiveKit..."); livekit::SignalClient client; - client.Connect("ws://localhost:7880/rtc", "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE2NTMzMTY2MjYsImlzcyI6IkFQSUNrSG04M01oZ2hQeCIsIm5iZiI6MTY1MDcyNDYyNiwic3ViIjoidGVzdCIsInZpZGVvIjp7InJvb20iOiJ0ZXN0cm9vbSIsInJvb21Kb2luIjp0cnVlfX0.I3Q5W4pk1kUNguGEJ4m95nE8hl8cPCliXBtF9hCt-Wg"); + client.Connect("ws://localhost:7880", "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE2NTMzMTY2MjYsImlzcyI6IkFQSUNrSG04M01oZ2hQeCIsIm5iZiI6MTY1MDcyNDYyNiwic3ViIjoidGVzdCIsInZpZGVvIjp7InJvb20iOiJ0ZXN0cm9vbSIsInJvb21Kb2luIjp0cnVlfX0.I3Q5W4pk1kUNguGEJ4m95nE8hl8cPCliXBtF9hCt-Wg"); - return 0; -} \ No newline at end of file + while(true){ + client.Update(); + } + + return EXIT_SUCCESS; +} diff --git a/protocol b/protocol index fa9efaf..377f74b 160000 --- a/protocol +++ b/protocol @@ -1 +1 @@ -Subproject commit fa9efaff8ca5cb1df5446f481afe73b8e89cda0f +Subproject commit 377f74b29d8d0fab17bb5778091144809f9c1123 diff --git a/src/proto/livekit_models.pb.cc b/src/proto/livekit_models.pb.cc index 750ce6f..307c6db 100644 --- a/src/proto/livekit_models.pb.cc +++ b/src/proto/livekit_models.pb.cc @@ -296,7 +296,9 @@ constexpr RTPStats::RTPStats( , rtt_current_(0u) , rtt_max_(0u) , key_frames_(0u) - , layer_lock_plis_(0u){} + , layer_lock_plis_(0u) + , nack_acks_(0u) + , nack_repeated_(0u){} struct RTPStatsDefaultTypeInternal { constexpr RTPStatsDefaultTypeInternal() : _instance(::PROTOBUF_NAMESPACE_ID::internal::ConstantInitialized{}) {} @@ -508,7 +510,9 @@ const uint32_t TableStruct_livekit_5fmodels_2eproto::offsets[] PROTOBUF_SECTION_ PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, jitter_max_), PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, gap_histogram_), PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, nacks_), + PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, nack_acks_), PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, nack_misses_), + PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, nack_repeated_), PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, plis_), PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, last_pli_), PROTOBUF_FIELD_OFFSET(::livekit::RTPStats, firs_), @@ -614,7 +618,7 @@ const char descriptor_table_protodef_livekit_5fmodels_2eproto[] PROTOBUF_SECTION "nnection\030\003 \001(\0162\034.livekit.ClientConfigSet" "ting\"L\n\022VideoConfiguration\0226\n\020hardware_e" "ncoder\030\001 \001(\0162\034.livekit.ClientConfigSetti" - "ng\"\237\010\n\010RTPStats\022.\n\nstart_time\030\001 \001(\0132\032.go" + "ng\"\311\010\n\010RTPStats\022.\n\nstart_time\030\001 \001(\0132\032.go" "ogle.protobuf.Timestamp\022,\n\010end_time\030\002 \001(" "\0132\032.google.protobuf.Timestamp\022\020\n\010duratio" "n\030\003 \001(\001\022\017\n\007packets\030\004 \001(\r\022\023\n\013packet_rate\030" @@ -630,34 +634,35 @@ const char descriptor_table_protodef_livekit_5fmodels_2eproto[] PROTOBUF_SECTION "\016\n\006frames\030\024 \001(\r\022\022\n\nframe_rate\030\025 \001(\001\022\026\n\016j" "itter_current\030\026 \001(\001\022\022\n\njitter_max\030\027 \001(\001\022" ":\n\rgap_histogram\030\030 \003(\0132#.livekit.RTPStat" - "s.GapHistogramEntry\022\r\n\005nacks\030\031 \001(\r\022\023\n\013na" - "ck_misses\030\032 \001(\r\022\014\n\004plis\030\033 \001(\r\022,\n\010last_pl" - "i\030\034 \001(\0132\032.google.protobuf.Timestamp\022\014\n\004f" - "irs\030\035 \001(\r\022,\n\010last_fir\030\036 \001(\0132\032.google.pro" - "tobuf.Timestamp\022\023\n\013rtt_current\030\037 \001(\r\022\017\n\007" - "rtt_max\030 \001(\r\022\022\n\nkey_frames\030! \001(\r\0222\n\016las" - "t_key_frame\030\" \001(\0132\032.google.protobuf.Time" - "stamp\022\027\n\017layer_lock_plis\030# \001(\r\0227\n\023last_l" - "ayer_lock_pli\030$ \001(\0132\032.google.protobuf.Ti" - "mestamp\0323\n\021GapHistogramEntry\022\013\n\003key\030\001 \001(" - "\005\022\r\n\005value\030\002 \001(\r:\0028\001*+\n\tTrackType\022\t\n\005AUD" - "IO\020\000\022\t\n\005VIDEO\020\001\022\010\n\004DATA\020\002*`\n\013TrackSource" - "\022\013\n\007UNKNOWN\020\000\022\n\n\006CAMERA\020\001\022\016\n\nMICROPHONE\020" - "\002\022\020\n\014SCREEN_SHARE\020\003\022\026\n\022SCREEN_SHARE_AUDI" - "O\020\004*6\n\014VideoQuality\022\007\n\003LOW\020\000\022\n\n\006MEDIUM\020\001" - "\022\010\n\004HIGH\020\002\022\007\n\003OFF\020\003*6\n\021ConnectionQuality" - "\022\010\n\004POOR\020\000\022\010\n\004GOOD\020\001\022\r\n\tEXCELLENT\020\002*;\n\023C" - "lientConfigSetting\022\t\n\005UNSET\020\000\022\014\n\010DISABLE" - "D\020\001\022\013\n\007ENABLED\020\002BFZ#github.com/livekit/p" - "rotocol/livekit\252\002\rLiveKit.Proto\352\002\016LiveKi" - "t::Protob\006proto3" + "s.GapHistogramEntry\022\r\n\005nacks\030\031 \001(\r\022\021\n\tna" + "ck_acks\030% \001(\r\022\023\n\013nack_misses\030\032 \001(\r\022\025\n\rna" + "ck_repeated\030& \001(\r\022\014\n\004plis\030\033 \001(\r\022,\n\010last_" + "pli\030\034 \001(\0132\032.google.protobuf.Timestamp\022\014\n" + "\004firs\030\035 \001(\r\022,\n\010last_fir\030\036 \001(\0132\032.google.p" + "rotobuf.Timestamp\022\023\n\013rtt_current\030\037 \001(\r\022\017" + "\n\007rtt_max\030 \001(\r\022\022\n\nkey_frames\030! \001(\r\0222\n\016l" + "ast_key_frame\030\" \001(\0132\032.google.protobuf.Ti" + "mestamp\022\027\n\017layer_lock_plis\030# \001(\r\0227\n\023last" + "_layer_lock_pli\030$ \001(\0132\032.google.protobuf." + "Timestamp\0323\n\021GapHistogramEntry\022\013\n\003key\030\001 " + "\001(\005\022\r\n\005value\030\002 \001(\r:\0028\001*+\n\tTrackType\022\t\n\005A" + "UDIO\020\000\022\t\n\005VIDEO\020\001\022\010\n\004DATA\020\002*`\n\013TrackSour" + "ce\022\013\n\007UNKNOWN\020\000\022\n\n\006CAMERA\020\001\022\016\n\nMICROPHON" + "E\020\002\022\020\n\014SCREEN_SHARE\020\003\022\026\n\022SCREEN_SHARE_AU" + "DIO\020\004*6\n\014VideoQuality\022\007\n\003LOW\020\000\022\n\n\006MEDIUM" + "\020\001\022\010\n\004HIGH\020\002\022\007\n\003OFF\020\003*6\n\021ConnectionQuali" + "ty\022\010\n\004POOR\020\000\022\010\n\004GOOD\020\001\022\r\n\tEXCELLENT\020\002*;\n" + "\023ClientConfigSetting\022\t\n\005UNSET\020\000\022\014\n\010DISAB" + "LED\020\001\022\013\n\007ENABLED\020\002BFZ#github.com/livekit" + "/protocol/livekit\252\002\rLiveKit.Proto\352\002\016Live" + "Kit::Protob\006proto3" ; static const ::PROTOBUF_NAMESPACE_ID::internal::DescriptorTable*const descriptor_table_livekit_5fmodels_2eproto_deps[1] = { &::descriptor_table_google_2fprotobuf_2ftimestamp_2eproto, }; static ::PROTOBUF_NAMESPACE_ID::internal::once_flag descriptor_table_livekit_5fmodels_2eproto_once; const ::PROTOBUF_NAMESPACE_ID::internal::DescriptorTable descriptor_table_livekit_5fmodels_2eproto = { - false, false, 3656, descriptor_table_protodef_livekit_5fmodels_2eproto, "livekit_models.proto", + false, false, 3698, descriptor_table_protodef_livekit_5fmodels_2eproto, "livekit_models.proto", &descriptor_table_livekit_5fmodels_2eproto_once, descriptor_table_livekit_5fmodels_2eproto_deps, 1, 16, schemas, file_default_instances, TableStruct_livekit_5fmodels_2eproto::offsets, file_level_metadata_livekit_5fmodels_2eproto, file_level_enum_descriptors_livekit_5fmodels_2eproto, file_level_service_descriptors_livekit_5fmodels_2eproto, @@ -5742,16 +5747,16 @@ RTPStats::RTPStats(const RTPStats& from) last_layer_lock_pli_ = nullptr; } ::memcpy(&duration_, &from.duration_, - static_cast(reinterpret_cast(&layer_lock_plis_) - - reinterpret_cast(&duration_)) + sizeof(layer_lock_plis_)); + static_cast(reinterpret_cast(&nack_repeated_) - + reinterpret_cast(&duration_)) + sizeof(nack_repeated_)); // @@protoc_insertion_point(copy_constructor:livekit.RTPStats) } inline void RTPStats::SharedCtor() { ::memset(reinterpret_cast(this) + static_cast( reinterpret_cast(&start_time_) - reinterpret_cast(this)), - 0, static_cast(reinterpret_cast(&layer_lock_plis_) - - reinterpret_cast(&start_time_)) + sizeof(layer_lock_plis_)); + 0, static_cast(reinterpret_cast(&nack_repeated_) - + reinterpret_cast(&start_time_)) + sizeof(nack_repeated_)); } RTPStats::~RTPStats() { @@ -5817,8 +5822,8 @@ void RTPStats::Clear() { } last_layer_lock_pli_ = nullptr; ::memset(&duration_, 0, static_cast( - reinterpret_cast(&layer_lock_plis_) - - reinterpret_cast(&duration_)) + sizeof(layer_lock_plis_)); + reinterpret_cast(&nack_repeated_) - + reinterpret_cast(&duration_)) + sizeof(nack_repeated_)); _internal_metadata_.Clear<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(); } @@ -6121,6 +6126,22 @@ const char* RTPStats::_InternalParse(const char* ptr, ::PROTOBUF_NAMESPACE_ID::i } else goto handle_unusual; continue; + // uint32 nack_acks = 37; + case 37: + if (PROTOBUF_PREDICT_TRUE(static_cast(tag) == 40)) { + nack_acks_ = ::PROTOBUF_NAMESPACE_ID::internal::ReadVarint32(&ptr); + CHK_(ptr); + } else + goto handle_unusual; + continue; + // uint32 nack_repeated = 38; + case 38: + if (PROTOBUF_PREDICT_TRUE(static_cast(tag) == 48)) { + nack_repeated_ = ::PROTOBUF_NAMESPACE_ID::internal::ReadVarint32(&ptr); + CHK_(ptr); + } else + goto handle_unusual; + continue; default: goto handle_unusual; } // switch @@ -6451,6 +6472,18 @@ uint8_t* RTPStats::_InternalSerialize( 36, _Internal::last_layer_lock_pli(this), target, stream); } + // uint32 nack_acks = 37; + if (this->_internal_nack_acks() != 0) { + target = stream->EnsureSpace(target); + target = ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::WriteUInt32ToArray(37, this->_internal_nack_acks(), target); + } + + // uint32 nack_repeated = 38; + if (this->_internal_nack_repeated() != 0) { + target = stream->EnsureSpace(target); + target = ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::WriteUInt32ToArray(38, this->_internal_nack_repeated(), target); + } + if (PROTOBUF_PREDICT_FALSE(_internal_metadata_.have_unknown_fields())) { target = ::PROTOBUF_NAMESPACE_ID::internal::WireFormat::InternalSerializeUnknownFieldsToArray( _internal_metadata_.unknown_fields<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(::PROTOBUF_NAMESPACE_ID::UnknownFieldSet::default_instance), target, stream); @@ -6733,6 +6766,20 @@ size_t RTPStats::ByteSizeLong() const { this->_internal_layer_lock_plis()); } + // uint32 nack_acks = 37; + if (this->_internal_nack_acks() != 0) { + total_size += 2 + + ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::UInt32Size( + this->_internal_nack_acks()); + } + + // uint32 nack_repeated = 38; + if (this->_internal_nack_repeated() != 0) { + total_size += 2 + + ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::UInt32Size( + this->_internal_nack_repeated()); + } + return MaybeComputeUnknownFieldsSize(total_size, &_cached_size_); } @@ -6909,6 +6956,12 @@ void RTPStats::MergeFrom(const RTPStats& from) { if (from._internal_layer_lock_plis() != 0) { _internal_set_layer_lock_plis(from._internal_layer_lock_plis()); } + if (from._internal_nack_acks() != 0) { + _internal_set_nack_acks(from._internal_nack_acks()); + } + if (from._internal_nack_repeated() != 0) { + _internal_set_nack_repeated(from._internal_nack_repeated()); + } _internal_metadata_.MergeFrom<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(from._internal_metadata_); } @@ -6928,8 +6981,8 @@ void RTPStats::InternalSwap(RTPStats* other) { _internal_metadata_.InternalSwap(&other->_internal_metadata_); gap_histogram_.InternalSwap(&other->gap_histogram_); ::PROTOBUF_NAMESPACE_ID::internal::memswap< - PROTOBUF_FIELD_OFFSET(RTPStats, layer_lock_plis_) - + sizeof(RTPStats::layer_lock_plis_) + PROTOBUF_FIELD_OFFSET(RTPStats, nack_repeated_) + + sizeof(RTPStats::nack_repeated_) - PROTOBUF_FIELD_OFFSET(RTPStats, start_time_)>( reinterpret_cast(&start_time_), reinterpret_cast(&other->start_time_)); diff --git a/src/proto/livekit_models.pb.h b/src/proto/livekit_models.pb.h index a463080..e48b778 100644 --- a/src/proto/livekit_models.pb.h +++ b/src/proto/livekit_models.pb.h @@ -3551,6 +3551,8 @@ class RTPStats final : kRttMaxFieldNumber = 32, kKeyFramesFieldNumber = 33, kLayerLockPlisFieldNumber = 35, + kNackAcksFieldNumber = 37, + kNackRepeatedFieldNumber = 38, }; // map gap_histogram = 24; int gap_histogram_size() const; @@ -3938,6 +3940,24 @@ class RTPStats final : void _internal_set_layer_lock_plis(uint32_t value); public: + // uint32 nack_acks = 37; + void clear_nack_acks(); + uint32_t nack_acks() const; + void set_nack_acks(uint32_t value); + private: + uint32_t _internal_nack_acks() const; + void _internal_set_nack_acks(uint32_t value); + public: + + // uint32 nack_repeated = 38; + void clear_nack_repeated(); + uint32_t nack_repeated() const; + void set_nack_repeated(uint32_t value); + private: + uint32_t _internal_nack_repeated() const; + void _internal_set_nack_repeated(uint32_t value); + public: + // @@protoc_insertion_point(class_scope:livekit.RTPStats) private: class _Internal; @@ -3985,6 +4005,8 @@ class RTPStats final : uint32_t rtt_max_; uint32_t key_frames_; uint32_t layer_lock_plis_; + uint32_t nack_acks_; + uint32_t nack_repeated_; mutable ::PROTOBUF_NAMESPACE_ID::internal::CachedSize _cached_size_; friend struct ::TableStruct_livekit_5fmodels_2eproto; }; @@ -7421,6 +7443,26 @@ inline void RTPStats::set_nacks(uint32_t value) { // @@protoc_insertion_point(field_set:livekit.RTPStats.nacks) } +// uint32 nack_acks = 37; +inline void RTPStats::clear_nack_acks() { + nack_acks_ = 0u; +} +inline uint32_t RTPStats::_internal_nack_acks() const { + return nack_acks_; +} +inline uint32_t RTPStats::nack_acks() const { + // @@protoc_insertion_point(field_get:livekit.RTPStats.nack_acks) + return _internal_nack_acks(); +} +inline void RTPStats::_internal_set_nack_acks(uint32_t value) { + + nack_acks_ = value; +} +inline void RTPStats::set_nack_acks(uint32_t value) { + _internal_set_nack_acks(value); + // @@protoc_insertion_point(field_set:livekit.RTPStats.nack_acks) +} + // uint32 nack_misses = 26; inline void RTPStats::clear_nack_misses() { nack_misses_ = 0u; @@ -7441,6 +7483,26 @@ inline void RTPStats::set_nack_misses(uint32_t value) { // @@protoc_insertion_point(field_set:livekit.RTPStats.nack_misses) } +// uint32 nack_repeated = 38; +inline void RTPStats::clear_nack_repeated() { + nack_repeated_ = 0u; +} +inline uint32_t RTPStats::_internal_nack_repeated() const { + return nack_repeated_; +} +inline uint32_t RTPStats::nack_repeated() const { + // @@protoc_insertion_point(field_get:livekit.RTPStats.nack_repeated) + return _internal_nack_repeated(); +} +inline void RTPStats::_internal_set_nack_repeated(uint32_t value) { + + nack_repeated_ = value; +} +inline void RTPStats::set_nack_repeated(uint32_t value) { + _internal_set_nack_repeated(value); + // @@protoc_insertion_point(field_set:livekit.RTPStats.nack_repeated) +} + // uint32 plis = 27; inline void RTPStats::clear_plis() { plis_ = 0u; diff --git a/src/proto/livekit_rtc.pb.cc b/src/proto/livekit_rtc.pb.cc index 86782ec..f6e58ac 100644 --- a/src/proto/livekit_rtc.pb.cc +++ b/src/proto/livekit_rtc.pb.cc @@ -339,6 +339,7 @@ constexpr TrackPermission::TrackPermission( ::PROTOBUF_NAMESPACE_ID::internal::ConstantInitialized) : track_sids_() , participant_sid_(&::PROTOBUF_NAMESPACE_ID::internal::fixed_address_empty_string) + , participant_identity_(&::PROTOBUF_NAMESPACE_ID::internal::fixed_address_empty_string) , all_tracks_(false){} struct TrackPermissionDefaultTypeInternal { constexpr TrackPermissionDefaultTypeInternal() @@ -657,6 +658,7 @@ const uint32_t TableStruct_livekit_5frtc_2eproto::offsets[] PROTOBUF_SECTION_VAR PROTOBUF_FIELD_OFFSET(::livekit::TrackPermission, participant_sid_), PROTOBUF_FIELD_OFFSET(::livekit::TrackPermission, all_tracks_), PROTOBUF_FIELD_OFFSET(::livekit::TrackPermission, track_sids_), + PROTOBUF_FIELD_OFFSET(::livekit::TrackPermission, participant_identity_), ~0u, // no _has_bits_ PROTOBUF_FIELD_OFFSET(::livekit::SubscriptionPermission, _internal_metadata_), ~0u, // no _extensions_ @@ -730,11 +732,11 @@ static const ::PROTOBUF_NAMESPACE_ID::internal::MigrationSchema schemas[] PROTOB { 208, -1, -1, sizeof(::livekit::SubscribedQuality)}, { 216, -1, -1, sizeof(::livekit::SubscribedQualityUpdate)}, { 224, -1, -1, sizeof(::livekit::TrackPermission)}, - { 233, -1, -1, sizeof(::livekit::SubscriptionPermission)}, - { 241, -1, -1, sizeof(::livekit::SubscriptionPermissionUpdate)}, - { 250, -1, -1, sizeof(::livekit::SyncState)}, - { 260, -1, -1, sizeof(::livekit::DataChannelInfo)}, - { 269, -1, -1, sizeof(::livekit::SimulateScenario)}, + { 234, -1, -1, sizeof(::livekit::SubscriptionPermission)}, + { 242, -1, -1, sizeof(::livekit::SubscriptionPermissionUpdate)}, + { 251, -1, -1, sizeof(::livekit::SyncState)}, + { 261, -1, -1, sizeof(::livekit::DataChannelInfo)}, + { 270, -1, -1, sizeof(::livekit::SimulateScenario)}, }; static ::PROTOBUF_NAMESPACE_ID::Message const * const file_default_instances[] = { @@ -858,36 +860,36 @@ const char descriptor_table_protodef_livekit_5frtc_2eproto[] PROTOBUF_SECTION_VA "eoQuality\022\017\n\007enabled\030\002 \001(\010\"f\n\027Subscribed" "QualityUpdate\022\021\n\ttrack_sid\030\001 \001(\t\0228\n\024subs" "cribed_qualities\030\002 \003(\0132\032.livekit.Subscri" - "bedQuality\"R\n\017TrackPermission\022\027\n\017partici" + "bedQuality\"p\n\017TrackPermission\022\027\n\017partici" "pant_sid\030\001 \001(\t\022\022\n\nall_tracks\030\002 \001(\010\022\022\n\ntr" - "ack_sids\030\003 \003(\t\"g\n\026SubscriptionPermission" - "\022\030\n\020all_participants\030\001 \001(\010\0223\n\021track_perm" - "issions\030\002 \003(\0132\030.livekit.TrackPermission\"" - "[\n\034SubscriptionPermissionUpdate\022\027\n\017parti" - "cipant_sid\030\001 \001(\t\022\021\n\ttrack_sid\030\002 \001(\t\022\017\n\007a" - "llowed\030\003 \001(\010\"\325\001\n\tSyncState\022+\n\006answer\030\001 \001" - "(\0132\033.livekit.SessionDescription\0221\n\014subsc" - "ription\030\002 \001(\0132\033.livekit.UpdateSubscripti" - "on\0227\n\016publish_tracks\030\003 \003(\0132\037.livekit.Tra" - "ckPublishedResponse\022/\n\rdata_channels\030\004 \003" - "(\0132\030.livekit.DataChannelInfo\"S\n\017DataChan" - "nelInfo\022\r\n\005label\030\001 \001(\t\022\n\n\002id\030\002 \001(\r\022%\n\006ta" - "rget\030\003 \001(\0162\025.livekit.SignalTarget\"}\n\020Sim" - "ulateScenario\022\030\n\016speaker_update\030\001 \001(\005H\000\022" - "\026\n\014node_failure\030\002 \001(\010H\000\022\023\n\tmigration\030\003 \001" - "(\010H\000\022\026\n\014server_leave\030\004 \001(\010H\000B\n\n\010scenario" - "*-\n\014SignalTarget\022\r\n\tPUBLISHER\020\000\022\016\n\nSUBSC" - "RIBER\020\001*%\n\013StreamState\022\n\n\006ACTIVE\020\000\022\n\n\006PA" - "USED\020\001BFZ#github.com/livekit/protocol/li" - "vekit\252\002\rLiveKit.Proto\352\002\016LiveKit::Protob\006" - "proto3" + "ack_sids\030\003 \003(\t\022\034\n\024participant_identity\030\004" + " \001(\t\"g\n\026SubscriptionPermission\022\030\n\020all_pa" + "rticipants\030\001 \001(\010\0223\n\021track_permissions\030\002 " + "\003(\0132\030.livekit.TrackPermission\"[\n\034Subscri" + "ptionPermissionUpdate\022\027\n\017participant_sid" + "\030\001 \001(\t\022\021\n\ttrack_sid\030\002 \001(\t\022\017\n\007allowed\030\003 \001" + "(\010\"\325\001\n\tSyncState\022+\n\006answer\030\001 \001(\0132\033.livek" + "it.SessionDescription\0221\n\014subscription\030\002 " + "\001(\0132\033.livekit.UpdateSubscription\0227\n\016publ" + "ish_tracks\030\003 \003(\0132\037.livekit.TrackPublishe" + "dResponse\022/\n\rdata_channels\030\004 \003(\0132\030.livek" + "it.DataChannelInfo\"S\n\017DataChannelInfo\022\r\n" + "\005label\030\001 \001(\t\022\n\n\002id\030\002 \001(\r\022%\n\006target\030\003 \001(\016" + "2\025.livekit.SignalTarget\"}\n\020SimulateScena" + "rio\022\030\n\016speaker_update\030\001 \001(\005H\000\022\026\n\014node_fa" + "ilure\030\002 \001(\010H\000\022\023\n\tmigration\030\003 \001(\010H\000\022\026\n\014se" + "rver_leave\030\004 \001(\010H\000B\n\n\010scenario*-\n\014Signal" + "Target\022\r\n\tPUBLISHER\020\000\022\016\n\nSUBSCRIBER\020\001*%\n" + "\013StreamState\022\n\n\006ACTIVE\020\000\022\n\n\006PAUSED\020\001BFZ#" + "github.com/livekit/protocol/livekit\252\002\rLi" + "veKit.Proto\352\002\016LiveKit::Protob\006proto3" ; static const ::PROTOBUF_NAMESPACE_ID::internal::DescriptorTable*const descriptor_table_livekit_5frtc_2eproto_deps[1] = { &::descriptor_table_livekit_5fmodels_2eproto, }; static ::PROTOBUF_NAMESPACE_ID::internal::once_flag descriptor_table_livekit_5frtc_2eproto_once; const ::PROTOBUF_NAMESPACE_ID::internal::DescriptorTable descriptor_table_livekit_5frtc_2eproto = { - false, false, 4406, descriptor_table_protodef_livekit_5frtc_2eproto, "livekit_rtc.proto", + false, false, 4436, descriptor_table_protodef_livekit_5frtc_2eproto, "livekit_rtc.proto", &descriptor_table_livekit_5frtc_2eproto_once, descriptor_table_livekit_5frtc_2eproto_deps, 1, 29, schemas, file_default_instances, TableStruct_livekit_5frtc_2eproto::offsets, file_level_metadata_livekit_5frtc_2eproto, file_level_enum_descriptors_livekit_5frtc_2eproto, file_level_service_descriptors_livekit_5frtc_2eproto, @@ -8274,6 +8276,14 @@ TrackPermission::TrackPermission(const TrackPermission& from) participant_sid_.Set(::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr::EmptyDefault{}, from._internal_participant_sid(), GetArenaForAllocation()); } + participant_identity_.UnsafeSetDefault(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited()); + #ifdef PROTOBUF_FORCE_COPY_DEFAULT_STRING + participant_identity_.Set(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), "", GetArenaForAllocation()); + #endif // PROTOBUF_FORCE_COPY_DEFAULT_STRING + if (!from._internal_participant_identity().empty()) { + participant_identity_.Set(::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr::EmptyDefault{}, from._internal_participant_identity(), + GetArenaForAllocation()); + } all_tracks_ = from.all_tracks_; // @@protoc_insertion_point(copy_constructor:livekit.TrackPermission) } @@ -8283,6 +8293,10 @@ participant_sid_.UnsafeSetDefault(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptySt #ifdef PROTOBUF_FORCE_COPY_DEFAULT_STRING participant_sid_.Set(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), "", GetArenaForAllocation()); #endif // PROTOBUF_FORCE_COPY_DEFAULT_STRING +participant_identity_.UnsafeSetDefault(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited()); +#ifdef PROTOBUF_FORCE_COPY_DEFAULT_STRING + participant_identity_.Set(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), "", GetArenaForAllocation()); +#endif // PROTOBUF_FORCE_COPY_DEFAULT_STRING all_tracks_ = false; } @@ -8296,6 +8310,7 @@ TrackPermission::~TrackPermission() { inline void TrackPermission::SharedDtor() { GOOGLE_DCHECK(GetArenaForAllocation() == nullptr); participant_sid_.DestroyNoArena(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited()); + participant_identity_.DestroyNoArena(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited()); } void TrackPermission::ArenaDtor(void* object) { @@ -8316,6 +8331,7 @@ void TrackPermission::Clear() { track_sids_.Clear(); participant_sid_.ClearToEmpty(); + participant_identity_.ClearToEmpty(); all_tracks_ = false; _internal_metadata_.Clear<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(); } @@ -8359,6 +8375,16 @@ const char* TrackPermission::_InternalParse(const char* ptr, ::PROTOBUF_NAMESPAC } else goto handle_unusual; continue; + // string participant_identity = 4; + case 4: + if (PROTOBUF_PREDICT_TRUE(static_cast(tag) == 34)) { + auto str = _internal_mutable_participant_identity(); + ptr = ::PROTOBUF_NAMESPACE_ID::internal::InlineGreedyStringParser(str, ptr, ctx); + CHK_(::PROTOBUF_NAMESPACE_ID::internal::VerifyUTF8(str, "livekit.TrackPermission.participant_identity")); + CHK_(ptr); + } else + goto handle_unusual; + continue; default: goto handle_unusual; } // switch @@ -8414,6 +8440,16 @@ uint8_t* TrackPermission::_InternalSerialize( target = stream->WriteString(3, s, target); } + // string participant_identity = 4; + if (!this->_internal_participant_identity().empty()) { + ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::VerifyUtf8String( + this->_internal_participant_identity().data(), static_cast(this->_internal_participant_identity().length()), + ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::SERIALIZE, + "livekit.TrackPermission.participant_identity"); + target = stream->WriteStringMaybeAliased( + 4, this->_internal_participant_identity(), target); + } + if (PROTOBUF_PREDICT_FALSE(_internal_metadata_.have_unknown_fields())) { target = ::PROTOBUF_NAMESPACE_ID::internal::WireFormat::InternalSerializeUnknownFieldsToArray( _internal_metadata_.unknown_fields<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(::PROTOBUF_NAMESPACE_ID::UnknownFieldSet::default_instance), target, stream); @@ -8445,6 +8481,13 @@ size_t TrackPermission::ByteSizeLong() const { this->_internal_participant_sid()); } + // string participant_identity = 4; + if (!this->_internal_participant_identity().empty()) { + total_size += 1 + + ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::StringSize( + this->_internal_participant_identity()); + } + // bool all_tracks = 2; if (this->_internal_all_tracks() != 0) { total_size += 1 + 1; @@ -8476,6 +8519,9 @@ void TrackPermission::MergeFrom(const TrackPermission& from) { if (!from._internal_participant_sid().empty()) { _internal_set_participant_sid(from._internal_participant_sid()); } + if (!from._internal_participant_identity().empty()) { + _internal_set_participant_identity(from._internal_participant_identity()); + } if (from._internal_all_tracks() != 0) { _internal_set_all_tracks(from._internal_all_tracks()); } @@ -8504,6 +8550,11 @@ void TrackPermission::InternalSwap(TrackPermission* other) { &participant_sid_, lhs_arena, &other->participant_sid_, rhs_arena ); + ::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr::InternalSwap( + &::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), + &participant_identity_, lhs_arena, + &other->participant_identity_, rhs_arena + ); swap(all_tracks_, other->all_tracks_); } diff --git a/src/proto/livekit_rtc.pb.h b/src/proto/livekit_rtc.pb.h index 6bac53f..265e8f0 100644 --- a/src/proto/livekit_rtc.pb.h +++ b/src/proto/livekit_rtc.pb.h @@ -4990,6 +4990,7 @@ class TrackPermission final : enum : int { kTrackSidsFieldNumber = 3, kParticipantSidFieldNumber = 1, + kParticipantIdentityFieldNumber = 4, kAllTracksFieldNumber = 2, }; // repeated string track_sids = 3; @@ -5030,6 +5031,20 @@ class TrackPermission final : std::string* _internal_mutable_participant_sid(); public: + // string participant_identity = 4; + void clear_participant_identity(); + const std::string& participant_identity() const; + template + void set_participant_identity(ArgT0&& arg0, ArgT... args); + std::string* mutable_participant_identity(); + PROTOBUF_NODISCARD std::string* release_participant_identity(); + void set_allocated_participant_identity(std::string* participant_identity); + private: + const std::string& _internal_participant_identity() const; + inline PROTOBUF_ALWAYS_INLINE void _internal_set_participant_identity(const std::string& value); + std::string* _internal_mutable_participant_identity(); + public: + // bool all_tracks = 2; void clear_all_tracks(); bool all_tracks() const; @@ -5048,6 +5063,7 @@ class TrackPermission final : typedef void DestructorSkippable_; ::PROTOBUF_NAMESPACE_ID::RepeatedPtrField track_sids_; ::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr participant_sid_; + ::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr participant_identity_; bool all_tracks_; mutable ::PROTOBUF_NAMESPACE_ID::internal::CachedSize _cached_size_; friend struct ::TableStruct_livekit_5frtc_2eproto; @@ -10704,6 +10720,57 @@ TrackPermission::mutable_track_sids() { return &track_sids_; } +// string participant_identity = 4; +inline void TrackPermission::clear_participant_identity() { + participant_identity_.ClearToEmpty(); +} +inline const std::string& TrackPermission::participant_identity() const { + // @@protoc_insertion_point(field_get:livekit.TrackPermission.participant_identity) + return _internal_participant_identity(); +} +template +inline PROTOBUF_ALWAYS_INLINE +void TrackPermission::set_participant_identity(ArgT0&& arg0, ArgT... args) { + + participant_identity_.Set(::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr::EmptyDefault{}, static_cast(arg0), args..., GetArenaForAllocation()); + // @@protoc_insertion_point(field_set:livekit.TrackPermission.participant_identity) +} +inline std::string* TrackPermission::mutable_participant_identity() { + std::string* _s = _internal_mutable_participant_identity(); + // @@protoc_insertion_point(field_mutable:livekit.TrackPermission.participant_identity) + return _s; +} +inline const std::string& TrackPermission::_internal_participant_identity() const { + return participant_identity_.Get(); +} +inline void TrackPermission::_internal_set_participant_identity(const std::string& value) { + + participant_identity_.Set(::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr::EmptyDefault{}, value, GetArenaForAllocation()); +} +inline std::string* TrackPermission::_internal_mutable_participant_identity() { + + return participant_identity_.Mutable(::PROTOBUF_NAMESPACE_ID::internal::ArenaStringPtr::EmptyDefault{}, GetArenaForAllocation()); +} +inline std::string* TrackPermission::release_participant_identity() { + // @@protoc_insertion_point(field_release:livekit.TrackPermission.participant_identity) + return participant_identity_.Release(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), GetArenaForAllocation()); +} +inline void TrackPermission::set_allocated_participant_identity(std::string* participant_identity) { + if (participant_identity != nullptr) { + + } else { + + } + participant_identity_.SetAllocated(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), participant_identity, + GetArenaForAllocation()); +#ifdef PROTOBUF_FORCE_COPY_DEFAULT_STRING + if (participant_identity_.IsDefault(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited())) { + participant_identity_.Set(&::PROTOBUF_NAMESPACE_ID::internal::GetEmptyStringAlreadyInited(), "", GetArenaForAllocation()); + } +#endif // PROTOBUF_FORCE_COPY_DEFAULT_STRING + // @@protoc_insertion_point(field_set_allocated:livekit.TrackPermission.participant_identity) +} + // ------------------------------------------------------------------- // SubscriptionPermission diff --git a/src/room.cpp b/src/room.cpp new file mode 100644 index 0000000..196d37d --- /dev/null +++ b/src/room.cpp @@ -0,0 +1,5 @@ +// +// Created by Théo Monnom on 07/05/2022. +// + +#include "room.h" diff --git a/src/room.h b/src/room.h new file mode 100644 index 0000000..2b305f9 --- /dev/null +++ b/src/room.h @@ -0,0 +1,27 @@ +// +// Created by Théo Monnom on 07/05/2022. +// + +#ifndef LIVEKIT_NATIVE_ROOM_H +#define LIVEKIT_NATIVE_ROOM_H + +#include +#include "signal_client.h" + +namespace livekit{ + /*class Room { + Room(); + ~Room(); + + void Connect(const std::string &url, const std::string &token); + + + + private: + SignalClient m_Client; + };*/ + +} + + +#endif //LIVEKIT_NATIVE_ROOM_H diff --git a/src/rtc_engine.cpp b/src/rtc_engine.cpp new file mode 100644 index 0000000..d654e16 --- /dev/null +++ b/src/rtc_engine.cpp @@ -0,0 +1,14 @@ +// +// Created by Théo Monnom on 04/05/2022. +// + +#include "rtc_engine.h" + +namespace livekit{ + + void RTCEngine::configure() { + + + } + +} diff --git a/src/rtc_engine.h b/src/rtc_engine.h new file mode 100644 index 0000000..37c7439 --- /dev/null +++ b/src/rtc_engine.h @@ -0,0 +1,25 @@ +// +// Created by Théo Monnom on 04/05/2022. +// + +#ifndef LIVEKIT_NATIVE_RTC_ENGINE_H +#define LIVEKIT_NATIVE_RTC_ENGINE_H + +#include "signal_client.h" + +namespace livekit{ + + class RTCEngine { + + + private: + void configure(); + + private: + SignalClient m_Client; + }; +}; + + + +#endif //LIVEKIT_NATIVE_RTC_ENGINE_H diff --git a/src/signal_client.cpp b/src/signal_client.cpp index 74e57f6..3cc5af4 100644 --- a/src/signal_client.cpp +++ b/src/signal_client.cpp @@ -3,36 +3,129 @@ // #include "signal_client.h" -#include "proto/livekit_rtc.pb.h" -#include -#include +#include +#include namespace livekit { - void SignalClient::Connect(const std::string& url, const std::string& token) { - boost::regex reg("(ws|wss)://([^:]+):?([^/]*)(.*)"); - boost::match_results regMatches; - if (boost::regex_match(url, regMatches, reg)) - { - const std::string protocol = regMatches[1]; // TODO Check if secure or not - const std::string domain = regMatches[2]; - const std::string port = regMatches[3]; - m_WebSocket = std::make_unique>(m_IOContext); + SignalClient::SignalClient() : m_Connected(false), m_Writing(false), m_Reading(false) { - tcp::resolver resolver{m_IOContext}; - auto const results = resolver.resolve(domain, port); - net::connect(m_WebSocket->next_layer(), results.begin(), results.end()); + } - m_WebSocket->handshake(domain, "/rtc?access_token=" + token); + SignalClient::~SignalClient() { + Disconnect(); + } - while(true){ - beast::flat_buffer buffer; - m_WebSocket->read(buffer); + void SignalClient::Connect(const std::string &url, const std::string &token) { + if (m_Connected) + throw std::runtime_error{"already connected"}; + + m_URL = ParseURL(url); + m_Token = token; + + Start(); // We don't need a thread, everything is async ( + easier to maintain ) + } + + void SignalClient::Update() { + beast::error_code ec; + m_IOContext.poll(ec); + + if (ec) + throw std::runtime_error{"SignalClient::Update - " + ec.message()}; + + if (m_Connected) { + if (!m_WebSocket.is_open()) + throw std::runtime_error{"Websocket isn't open"}; // TODO Start reconnect + + if (!m_Reading) { + m_WebSocket.async_read(m_ReadBuffer, beast::bind_front_handler(&SignalClient::OnRead, this)); + m_Reading = true; } - }else{ - throw std::runtime_error{"Failed to parse url"}; + // Write pending messages + if (!m_Writing && !m_WriteQueue.empty()) { + auto req = m_WriteQueue.front(); + + unsigned long len = req.ByteSizeLong(); + uint8_t data[len]; + req.SerializeToArray(data, len); + + m_WebSocket.async_write(net::buffer(&data, len), + beast::bind_front_handler(&SignalClient::OnWrite, this)); + + m_WriteQueue.pop(); + m_Writing = true; + } } } + void SignalClient::Start() { + m_Resolver.async_resolve(m_URL.host, m_URL.port, beast::bind_front_handler(&SignalClient::OnResolve, this)); + } + + void SignalClient::Disconnect() { + if (!m_Connected) + return; + + m_Connected = false; + m_Work.reset(); + //m_IOContext.stop(); + m_WebSocket.close(websocket::close_code::normal); // TODO Close should be async + } + + void SignalClient::Send(SignalRequest req) { + m_WriteQueue.emplace(req); + } + + void SignalClient::OnResolve(beast::error_code ec, tcp::resolver::results_type results) { + if (ec) + throw std::runtime_error{"SignalClient::OnResolve - " + ec.message()}; + + auto &layer = beast::get_lowest_layer(m_WebSocket); + layer.expires_after(std::chrono::seconds(15)); + layer.async_connect(results, beast::bind_front_handler(&SignalClient::OnConnect, this)); + } + + void SignalClient::OnConnect(beast::error_code ec, tcp::resolver::results_type::endpoint_type ep) { + if (ec) + throw std::runtime_error{"SignalClient::OnConnect - " + ec.message()}; + + beast::get_lowest_layer(m_WebSocket).expires_never(); + m_WebSocket.set_option(websocket::stream_base::timeout::suggested(beast::role_type::client)); + + m_WebSocket.async_handshake(m_URL.host, "/rtc?access_token=" + m_Token + "&protocol=7", + beast::bind_front_handler(&SignalClient::OnHandshake, this)); + } + + void SignalClient::OnHandshake(beast::error_code ec) { + if (ec) + throw std::runtime_error{ + "SignalClient::OnHandshake - " + ec.message()}; // TODO Callback for handling errors + + m_Connected = true; + spdlog::info("Connected to Websocket"); + } + + void SignalClient::OnRead(beast::error_code ec, std::size_t bytesTransferred) { + m_Reading = false; + + if (ec) + throw std::runtime_error{"SignalClient::OnRead - " + ec.message()}; + + SignalResponse res{}; + if (res.ParseFromArray(m_ReadBuffer.cdata().data(), bytesTransferred)) { + m_ReadQueue.emplace(res); + } else { + spdlog::error("Failed to decode signal message"); + } + + m_ReadBuffer.clear(); + } + + void SignalClient::OnWrite(beast::error_code ec, std::size_t bytesTransferred) { + m_Writing = false; + + if (ec) + throw std::runtime_error{"SignalClient::OnWrite - " + ec.message()}; + } } \ No newline at end of file diff --git a/src/signal_client.h b/src/signal_client.h index 234b2ca..344e542 100644 --- a/src/signal_client.h +++ b/src/signal_client.h @@ -7,6 +7,9 @@ #include #include +#include +#include "proto/livekit_rtc.pb.h" +#include "utils.h" namespace beast = boost::beast; // from namespace http = beast::http; // from @@ -18,11 +21,43 @@ namespace livekit { class SignalClient { public: - void Connect(const std::string& url, const std::string& token); + SignalClient(); + ~SignalClient(); + + void Connect(const std::string &url, const std::string &token); + void Disconnect(); + void Update(); + void Send(SignalRequest req); + SignalResponse Poll(); private: + void Start(); + + // beast handlers + void OnResolve(beast::error_code ec, tcp::resolver::results_type results); + void OnConnect(beast::error_code ec, tcp::resolver::results_type::endpoint_type ep); + void OnHandshake(beast::error_code ec); + void OnRead(beast::error_code ec, std::size_t bytesTransferred); + void OnWrite(beast::error_code ec, std::size_t bytesTransferred); + + private: + URL m_URL; + std::string m_Token; + + std::queue m_ReadQueue; + std::queue m_WriteQueue; + bool m_Connected; + bool m_Reading, m_Writing; + + beast::flat_buffer m_WriteBuffer; + beast::flat_buffer m_ReadBuffer; + + // Keep order net::io_context m_IOContext; - std::unique_ptr> m_WebSocket; + net::executor_work_guard m_Work = net::make_work_guard( + m_IOContext); // Prevent the IOContext from running out of work + tcp::resolver m_Resolver{m_IOContext}; + websocket::stream m_WebSocket{m_IOContext}; }; } diff --git a/src/utils.h b/src/utils.h new file mode 100644 index 0000000..51762b3 --- /dev/null +++ b/src/utils.h @@ -0,0 +1,36 @@ +// +// Created by Théo Monnom on 01/05/2022. +// + +#ifndef LIVEKIT_NATIVE_UTILS_H +#define LIVEKIT_NATIVE_UTILS_H + +#include +#include + +namespace livekit { + + // TODO Do I need path + query ? + struct URL { + std::string protocol; + std::string host; + std::string port; + }; + + static URL ParseURL(const std::string &url) { + boost::regex reg("(ws|wss)://([^:]*):?(\\d*)(.*)"); + boost::match_results groups; + + if (boost::regex_search(url, groups, reg)) { + return URL{ + groups[1], + groups[2], + groups[3] + }; + } + + throw std::runtime_error{"failed to parse url"}; + } +} + +#endif //LIVEKIT_NATIVE_UTILS_H