Skip to content

Commit 99c8903

Browse files
committed
Fix tests
1 parent 8fc59cd commit 99c8903

18 files changed

Lines changed: 190 additions & 262 deletions

src/internal_modules/roc_netio/target_libuv/roc_netio/network_loop.cpp

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,11 @@ NetworkLoop::NetworkLoop(core::IPool& packet_pool,
140140
task_sem_.data = this;
141141
task_sem_initialized_ = true;
142142

143-
enable_realtime();
143+
if (!enable_realtime()) {
144+
roc_log(LogInfo,
145+
"network loop: can't set realtime priority of network thread. May need "
146+
"to be root");
147+
}
144148
if (!(started_ = Thread::start())) {
145149
init_status_ = status::StatusErrThread;
146150
return;

src/internal_modules/roc_node/receiver_decoder.cpp

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -169,6 +169,8 @@ status::StatusCode ReceiverDecoder::write_packet(address::Interface iface,
169169
roc_panic_if(!bytes);
170170
roc_panic_if(n_bytes == 0);
171171

172+
const core::nanoseconds_t capture_ts = core::timestamp(core::ClockUnix);
173+
172174
if (n_bytes > packet_factory_.packet_buffer_size()) {
173175
roc_log(LogError,
174176
"receiver decoder node:"
@@ -195,6 +197,7 @@ status::StatusCode ReceiverDecoder::write_packet(address::Interface iface,
195197
}
196198

197199
packet->add_flags(packet::Packet::FlagUDP);
200+
packet->udp()->receive_timestamp = capture_ts;
198201
packet->set_buffer(buffer);
199202

200203
packet::IWriter* writer = endpoint_writers_[iface];

src/internal_modules/roc_node/sender_encoder.cpp

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -226,6 +226,8 @@ SenderEncoder::write_packet(address::Interface iface, const void* bytes, size_t
226226
roc_panic_if(!bytes);
227227
roc_panic_if(n_bytes == 0);
228228

229+
const core::nanoseconds_t capture_ts = core::timestamp(core::ClockUnix);
230+
229231
if (n_bytes > packet_factory_.packet_buffer_size()) {
230232
roc_log(LogError,
231233
"sender encoder node:"
@@ -252,6 +254,7 @@ SenderEncoder::write_packet(address::Interface iface, const void* bytes, size_t
252254
}
253255

254256
packet->add_flags(packet::Packet::FlagUDP);
257+
packet->udp()->receive_timestamp = capture_ts;
255258
packet->set_buffer(buffer);
256259

257260
packet::IWriter* writer = endpoint_writers_[iface];

src/internal_modules/roc_pipeline/receiver_endpoint.cpp

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -199,7 +199,7 @@ packet::IWriter& ReceiverEndpoint::inbound_writer() {
199199
return *this;
200200
}
201201

202-
status::StatusCode ReceiverEndpoint::pull_packets(core::nanoseconds_t current_time) {
202+
status::StatusCode ReceiverEndpoint::pull_packets() {
203203
roc_panic_if(init_status_ != status::StatusOK);
204204

205205
roc_panic_if(!parser_);
@@ -209,7 +209,7 @@ status::StatusCode ReceiverEndpoint::pull_packets(core::nanoseconds_t current_ti
209209
// queue were added in a very short time or are being added currently. It's
210210
// acceptable to consider such packets late and pull them next time.
211211
while (packet::PacketPtr packet = inbound_queue_.try_pop_front_exclusive()) {
212-
const status::StatusCode code = handle_packet_(packet, current_time);
212+
const status::StatusCode code = handle_packet_(packet);
213213
state_tracker_.unregister_packet();
214214

215215
if (code != status::StatusOK) {
@@ -220,14 +220,13 @@ status::StatusCode ReceiverEndpoint::pull_packets(core::nanoseconds_t current_ti
220220
return status::StatusOK;
221221
}
222222

223-
status::StatusCode ReceiverEndpoint::handle_packet_(const packet::PacketPtr& packet,
224-
core::nanoseconds_t current_time) {
223+
status::StatusCode ReceiverEndpoint::handle_packet_(const packet::PacketPtr& packet) {
225224
if (!parser_->parse(*packet, packet->buffer())) {
226225
roc_log(LogDebug, "receiver endpoint: dropping bad packet: can't parse");
227226
return status::StatusOK;
228227
}
229228

230-
const status::StatusCode code = session_group_.route_packet(packet, current_time);
229+
const status::StatusCode code = session_group_.route_packet(packet);
231230

232231
if (code == status::StatusNoRoute) {
233232
roc_log(LogDebug, "receiver endpoint: dropping bad packet: can't route");

src/internal_modules/roc_pipeline/receiver_endpoint.h

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -93,13 +93,12 @@ class ReceiverEndpoint : public core::RefCounted<ReceiverEndpoint, core::ArenaAl
9393
//! Packets are written to inbound_writer() from network thread.
9494
//! They don't appear in pipeline immediately. Instead, pipeline thread
9595
//! should periodically call pull_packets() to make them available.
96-
ROC_ATTR_NODISCARD status::StatusCode pull_packets(core::nanoseconds_t current_time);
96+
ROC_ATTR_NODISCARD status::StatusCode pull_packets();
9797

9898
private:
9999
virtual ROC_ATTR_NODISCARD status::StatusCode write(const packet::PacketPtr& packet);
100100

101-
status::StatusCode handle_packet_(const packet::PacketPtr& packet,
102-
core::nanoseconds_t current_time);
101+
status::StatusCode handle_packet_(const packet::PacketPtr& packet);
103102

104103
const address::Protocol proto_;
105104

src/internal_modules/roc_pipeline/receiver_session_group.cpp

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -140,12 +140,11 @@ void ReceiverSessionGroup::reclock_sessions(core::nanoseconds_t playback_time) {
140140
}
141141
}
142142

143-
status::StatusCode ReceiverSessionGroup::route_packet(const packet::PacketPtr& packet,
144-
core::nanoseconds_t current_time) {
143+
status::StatusCode ReceiverSessionGroup::route_packet(const packet::PacketPtr& packet) {
145144
roc_panic_if(init_status_ != status::StatusOK);
146145

147146
if (packet->has_flags(packet::Packet::FlagControl)) {
148-
return route_control_packet_(packet, current_time);
147+
return route_control_packet_(packet);
149148
}
150149

151150
return route_transport_packet_(packet);
@@ -344,15 +343,14 @@ ReceiverSessionGroup::route_transport_packet_(const packet::PacketPtr& packet) {
344343
}
345344

346345
status::StatusCode
347-
ReceiverSessionGroup::route_control_packet_(const packet::PacketPtr& packet,
348-
core::nanoseconds_t current_time) {
346+
ReceiverSessionGroup::route_control_packet_(const packet::PacketPtr& packet) {
349347
if (!rtcp_communicator_) {
350348
roc_panic("session group: rtcp communicator is null");
351349
}
352350

353351
// This will invoke IParticipant methods implemented by us,
354352
// in particular notify_recv_stream() and maybe halt_recv_stream().
355-
return rtcp_communicator_->process_packet(packet, current_time);
353+
return rtcp_communicator_->process_packet(packet);
356354
}
357355

358356
bool ReceiverSessionGroup::can_create_session_(const packet::PacketPtr& packet) {

src/internal_modules/roc_pipeline/receiver_session_group.h

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -91,8 +91,7 @@ class ReceiverSessionGroup : public core::NonCopyable<>, private rtcp::IParticip
9191
void reclock_sessions(core::nanoseconds_t playback_time);
9292

9393
//! Route packet to session.
94-
ROC_ATTR_NODISCARD status::StatusCode route_packet(const packet::PacketPtr& packet,
95-
core::nanoseconds_t current_time);
94+
ROC_ATTR_NODISCARD status::StatusCode route_packet(const packet::PacketPtr& packet);
9695

9796
//! Get number of sessions in group.
9897
size_t num_sessions() const;
@@ -130,8 +129,7 @@ class ReceiverSessionGroup : public core::NonCopyable<>, private rtcp::IParticip
130129
virtual void halt_recv_stream(packet::stream_source_t send_source_id);
131130

132131
status::StatusCode route_transport_packet_(const packet::PacketPtr& packet);
133-
status::StatusCode route_control_packet_(const packet::PacketPtr& packet,
134-
core::nanoseconds_t current_time);
132+
status::StatusCode route_control_packet_(const packet::PacketPtr& packet);
135133

136134
bool can_create_session_(const packet::PacketPtr& packet);
137135

src/internal_modules/roc_pipeline/receiver_slot.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -84,19 +84,19 @@ status::StatusCode ReceiverSlot::refresh(core::nanoseconds_t current_time,
8484
status::StatusCode code = status::NoStatus;
8585

8686
if (source_endpoint_) {
87-
if ((code = source_endpoint_->pull_packets(current_time)) != status::StatusOK) {
87+
if ((code = source_endpoint_->pull_packets()) != status::StatusOK) {
8888
return code;
8989
}
9090
}
9191

9292
if (repair_endpoint_) {
93-
if ((code = repair_endpoint_->pull_packets(current_time)) != status::StatusOK) {
93+
if ((code = repair_endpoint_->pull_packets()) != status::StatusOK) {
9494
return code;
9595
}
9696
}
9797

9898
if (control_endpoint_) {
99-
if ((code = control_endpoint_->pull_packets(current_time)) != status::StatusOK) {
99+
if ((code = control_endpoint_->pull_packets()) != status::StatusOK) {
100100
return code;
101101
}
102102
}

src/internal_modules/roc_pipeline/sender_endpoint.cpp

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -184,7 +184,7 @@ packet::IWriter* SenderEndpoint::inbound_writer() {
184184
return this;
185185
}
186186

187-
status::StatusCode SenderEndpoint::pull_packets(core::nanoseconds_t current_time) {
187+
status::StatusCode SenderEndpoint::pull_packets() {
188188
roc_panic_if(init_status_ != status::StatusOK);
189189

190190
if (!parser_) {
@@ -197,7 +197,7 @@ status::StatusCode SenderEndpoint::pull_packets(core::nanoseconds_t current_time
197197
// queue were added in a very short time or are being added currently. It's
198198
// acceptable to consider such packets late and pull them next time.
199199
while (packet::PacketPtr packet = inbound_queue_.try_pop_front_exclusive()) {
200-
const status::StatusCode code = handle_packet_(packet, current_time);
200+
const status::StatusCode code = handle_packet_(packet);
201201
state_tracker_.unregister_packet();
202202

203203
if (code != status::StatusOK) {
@@ -208,14 +208,13 @@ status::StatusCode SenderEndpoint::pull_packets(core::nanoseconds_t current_time
208208
return status::StatusOK;
209209
}
210210

211-
status::StatusCode SenderEndpoint::handle_packet_(const packet::PacketPtr& packet,
212-
core::nanoseconds_t current_time) {
211+
status::StatusCode SenderEndpoint::handle_packet_(const packet::PacketPtr& packet) {
213212
if (!parser_->parse(*packet, packet->buffer())) {
214213
roc_log(LogDebug, "sender endpoint: dropping bad packet: can't parse");
215214
return status::StatusOK;
216215
}
217216

218-
const status::StatusCode code = sender_session_.route_packet(packet, current_time);
217+
const status::StatusCode code = sender_session_.route_packet(packet);
219218

220219
if (code == status::StatusNoRoute) {
221220
roc_log(LogDebug, "sender endpoint: dropping bad packet: can't route");

src/internal_modules/roc_pipeline/sender_endpoint.h

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -91,13 +91,12 @@ class SenderEndpoint : public core::NonCopyable<>, private packet::IWriter {
9191
//! Packets are written to inbound_writer() from network thread.
9292
//! They don't appear in pipeline immediately. Instead, pipeline thread
9393
//! should periodically call pull_packets() to make them available.
94-
ROC_ATTR_NODISCARD status::StatusCode pull_packets(core::nanoseconds_t current_time);
94+
ROC_ATTR_NODISCARD status::StatusCode pull_packets();
9595

9696
private:
9797
virtual ROC_ATTR_NODISCARD status::StatusCode write(const packet::PacketPtr& packet);
9898

99-
status::StatusCode handle_packet_(const packet::PacketPtr& packet,
100-
core::nanoseconds_t current_time);
99+
status::StatusCode handle_packet_(const packet::PacketPtr& packet);
101100

102101
const address::Protocol proto_;
103102

0 commit comments

Comments
 (0)