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
6 changes: 5 additions & 1 deletion src/c/core/session/session.c
Original file line number Diff line number Diff line change
Expand Up @@ -629,7 +629,11 @@ void uxr_flash_output_streams(
while (uxr_prepare_next_reliable_buffer_to_send(stream, &buffer, &length, &seq_num))
{
uxr_stamp_session_header(&session->info, id.raw, seq_num, buffer);
send_message(session, buffer, length);
if (!send_message(session, buffer, length))
{
uxr_cancel_reliable_buffer_to_send(stream, seq_num);
break;
}
}

UXR_UNLOCK_STREAM_ID(session, id);
Expand Down
13 changes: 13 additions & 0 deletions src/c/core/session/stream/output_reliable_stream.c
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,19 @@ bool uxr_prepare_next_reliable_buffer_to_send(
return data_to_send;
}

void uxr_cancel_reliable_buffer_to_send(
uxrOutputReliableStream* stream,
uxrSeqNum seq_num)
{
uxrSeqNum next_seq_num = uxr_seq_num_add(seq_num, 1);
stream->last_sent = uxr_seq_num_sub(seq_num, 1);
if (stream->last_written == next_seq_num &&
stream->offset == uxr_get_reliable_buffer_size(&stream->base, next_seq_num))
{
stream->last_written = seq_num;
}
}

bool uxr_update_output_stream_heartbeat_timestamp(
uxrOutputReliableStream* stream,
int64_t current_timestamp)
Expand Down
3 changes: 3 additions & 0 deletions src/c/core/session/stream/output_reliable_stream_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,9 @@ bool uxr_prepare_next_reliable_buffer_to_send(
uint8_t** buffer,
size_t* length,
uxrSeqNum* seq_num);
void uxr_cancel_reliable_buffer_to_send(
uxrOutputReliableStream* stream,
uxrSeqNum seq_num);

bool uxr_update_output_stream_heartbeat_timestamp(
uxrOutputReliableStream* stream,
Expand Down
31 changes: 31 additions & 0 deletions test/unitary/session/Session.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,7 @@ class SessionTest : public testing::Test
uint8_t input_reliable_buffer[MTU * HISTORY];

static int listening_counter;
static int sending_counter;

static bool send_msg(
void* instance,
Expand All @@ -124,6 +125,13 @@ class SessionTest : public testing::Test
EXPECT_EQ(size_t(MTU), len);
return false;
}
else if (std::string("FlashReliableStreamSendError") ==
::testing::UnitTest::GetInstance()->current_test_info()->name())
{
EXPECT_EQ(size_t(OFFSET + SUBHEADER_SIZE + 8), len);
SessionTest::sending_counter++;
return false;
}
else if (std::string("SendHeartbeat") == ::testing::UnitTest::GetInstance()->current_test_info()->name())
{
EXPECT_EQ(size_t(HEARTBEAT_MAX_MSG_SIZE - (MAX_HEADER_SIZE - OFFSET)), len);
Expand Down Expand Up @@ -272,6 +280,7 @@ class SessionTest : public testing::Test

SessionTest* SessionTest::current = nullptr;
int SessionTest::listening_counter;
int SessionTest::sending_counter;

TEST_F(SessionTest, SetStatusCallback)
{
Expand Down Expand Up @@ -355,6 +364,28 @@ TEST_F(SessionTest, FlashStreams)
uxr_flash_output_streams(&session);
}

TEST_F(SessionTest, FlashReliableStreamSendError)
{
SessionTest::sending_counter = 0;

ucdrBuffer ub;
uxrStreamId output_reliable = uxr_stream_id(0, UXR_RELIABLE_STREAM, UXR_OUTPUT_STREAM);
(void) uxr_prepare_stream_to_write_submessage(&session, output_reliable, 8, &ub, 1, 0);

uxrOutputReliableStream* stream = &session.streams.output_reliable[0];
uxr_flash_output_streams(&session);

EXPECT_EQ(1, SessionTest::sending_counter);
EXPECT_EQ(SEQ_NUM_MAX, stream->last_sent);
EXPECT_EQ(0u, stream->last_written);

uxr_flash_output_streams(&session);

EXPECT_EQ(2, SessionTest::sending_counter);
EXPECT_EQ(SEQ_NUM_MAX, stream->last_sent);
EXPECT_EQ(0u, stream->last_written);
}

TEST_F(SessionTest, WaitSessionStatusBad)
{
// The OK version is already checked with the CreateOk and DeleteOk test versions
Expand Down
Loading