Skip to content

Commit ebaf891

Browse files
committed
Fix reliable stream send failure handling
Generated-by: OpenAI Codex Signed-off-by: aineoae86-sys <ai.neo.ae86@gmail.com>
1 parent 0840e72 commit ebaf891

4 files changed

Lines changed: 52 additions & 1 deletion

File tree

src/c/core/session/session.c

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -629,7 +629,11 @@ void uxr_flash_output_streams(
629629
while (uxr_prepare_next_reliable_buffer_to_send(stream, &buffer, &length, &seq_num))
630630
{
631631
uxr_stamp_session_header(&session->info, id.raw, seq_num, buffer);
632-
send_message(session, buffer, length);
632+
if (!send_message(session, buffer, length))
633+
{
634+
uxr_cancel_reliable_buffer_to_send(stream, seq_num);
635+
break;
636+
}
633637
}
634638

635639
UXR_UNLOCK_STREAM_ID(session, id);

src/c/core/session/stream/output_reliable_stream.c

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -192,6 +192,19 @@ bool uxr_prepare_next_reliable_buffer_to_send(
192192
return data_to_send;
193193
}
194194

195+
void uxr_cancel_reliable_buffer_to_send(
196+
uxrOutputReliableStream* stream,
197+
uxrSeqNum seq_num)
198+
{
199+
uxrSeqNum next_seq_num = uxr_seq_num_add(seq_num, 1);
200+
stream->last_sent = uxr_seq_num_sub(seq_num, 1);
201+
if (stream->last_written == next_seq_num &&
202+
stream->offset == uxr_get_reliable_buffer_size(&stream->base, next_seq_num))
203+
{
204+
stream->last_written = seq_num;
205+
}
206+
}
207+
195208
bool uxr_update_output_stream_heartbeat_timestamp(
196209
uxrOutputReliableStream* stream,
197210
int64_t current_timestamp)

src/c/core/session/stream/output_reliable_stream_internal.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,9 @@ bool uxr_prepare_next_reliable_buffer_to_send(
4545
uint8_t** buffer,
4646
size_t* length,
4747
uxrSeqNum* seq_num);
48+
void uxr_cancel_reliable_buffer_to_send(
49+
uxrOutputReliableStream* stream,
50+
uxrSeqNum seq_num);
4851

4952
bool uxr_update_output_stream_heartbeat_timestamp(
5053
uxrOutputReliableStream* stream,

test/unitary/session/Session.cpp

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,7 @@ class SessionTest : public testing::Test
103103
uint8_t input_reliable_buffer[MTU * HISTORY];
104104

105105
static int listening_counter;
106+
static int sending_counter;
106107

107108
static bool send_msg(
108109
void* instance,
@@ -124,6 +125,13 @@ class SessionTest : public testing::Test
124125
EXPECT_EQ(size_t(MTU), len);
125126
return false;
126127
}
128+
else if (std::string("FlashReliableStreamSendError") ==
129+
::testing::UnitTest::GetInstance()->current_test_info()->name())
130+
{
131+
EXPECT_EQ(size_t(OFFSET + SUBHEADER_SIZE + 8), len);
132+
SessionTest::sending_counter++;
133+
return false;
134+
}
127135
else if (std::string("SendHeartbeat") == ::testing::UnitTest::GetInstance()->current_test_info()->name())
128136
{
129137
EXPECT_EQ(size_t(HEARTBEAT_MAX_MSG_SIZE - (MAX_HEADER_SIZE - OFFSET)), len);
@@ -272,6 +280,7 @@ class SessionTest : public testing::Test
272280

273281
SessionTest* SessionTest::current = nullptr;
274282
int SessionTest::listening_counter;
283+
int SessionTest::sending_counter;
275284

276285
TEST_F(SessionTest, SetStatusCallback)
277286
{
@@ -355,6 +364,28 @@ TEST_F(SessionTest, FlashStreams)
355364
uxr_flash_output_streams(&session);
356365
}
357366

367+
TEST_F(SessionTest, FlashReliableStreamSendError)
368+
{
369+
SessionTest::sending_counter = 0;
370+
371+
ucdrBuffer ub;
372+
uxrStreamId output_reliable = uxr_stream_id(0, UXR_RELIABLE_STREAM, UXR_OUTPUT_STREAM);
373+
(void) uxr_prepare_stream_to_write_submessage(&session, output_reliable, 8, &ub, 1, 0);
374+
375+
uxrOutputReliableStream* stream = &session.streams.output_reliable[0];
376+
uxr_flash_output_streams(&session);
377+
378+
EXPECT_EQ(1, SessionTest::sending_counter);
379+
EXPECT_EQ(SEQ_NUM_MAX, stream->last_sent);
380+
EXPECT_EQ(0u, stream->last_written);
381+
382+
uxr_flash_output_streams(&session);
383+
384+
EXPECT_EQ(2, SessionTest::sending_counter);
385+
EXPECT_EQ(SEQ_NUM_MAX, stream->last_sent);
386+
EXPECT_EQ(0u, stream->last_written);
387+
}
388+
358389
TEST_F(SessionTest, WaitSessionStatusBad)
359390
{
360391
// The OK version is already checked with the CreateOk and DeleteOk test versions

0 commit comments

Comments
 (0)