Skip to content

Commit eeb83d4

Browse files
RMW QoS incoming buffer handling (backport #165) (#176)
* RMW QoS incoming buffer handling (#165) * Initial * Update * Uncrustify * Fixes * Tests pubsub * Add two subs test * Uncrustify * Initial reqres * Adding complete req res tests * Uncrustify * Use absolute time reference * Avoid future data * Speedup tests * Update * Update * Pub/sub in best effort * Fix * Revert "Pub/sub in best effort" This reverts commit 9f46d80. * Revert "Fix" This reverts commit cb7a130. * Major review of rmw_wait * Uncrustify * Adapt Cmake * client branch * Update * Update * Uncrustify * Update rmw_microxrcedds_c/CMakeLists.txt Co-authored-by: Antonio Cuadros <49162117+Acuadros95@users.noreply.github.com> * Update rmw_microxrcedds_c/src/config.h.in Co-authored-by: Antonio Cuadros <49162117+Acuadros95@users.noreply.github.com> * Fix on_topic callback * Fix timeout handling * Exit if deserialize fails * Update .github/workflows/ci.yml Co-authored-by: Pablo Garrido <pablogs9@gmail.com> Co-authored-by: Antonio Cuadros <49162117+Acuadros95@users.noreply.github.com> Co-authored-by: Antonio cuadros <acuadros1995@gmail.com> (cherry picked from commit b1531c7) # Conflicts: # rmw_microxrcedds_c/src/types.h * Fix conflict * Update Co-authored-by: Pablo Garrido <pablogs9@gmail.com>
1 parent d6dcaeb commit eeb83d4

15 files changed

Lines changed: 1067 additions & 418 deletions

rmw_microxrcedds_c/CMakeLists.txt

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -121,10 +121,6 @@ if(${RMW_UXRCE_CREATION_MODE} STREQUAL "refs")
121121
set(RMW_UXRCE_USE_REFS ON)
122122
endif()
123123

124-
# Create source files with the define
125-
configure_file(${PROJECT_SOURCE_DIR}/src/config.h.in
126-
${PROJECT_BINARY_DIR}/include/rmw_microxrcedds_c/config.h)
127-
128124
# Set install directories
129125
if(WIN32)
130126
set(DOC_DIR "doc")
@@ -293,12 +289,19 @@ if(BUILD_TESTING)
293289
# Pedantic in CI
294290
set(CMAKE_C_FLAGS "${CMAKE_C_FLAGS} -Wall -Werror")
295291
set(CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS} -Wall -Werror")
292+
if(RMW_UXRCE_MAX_SESSIONS STREQUAL "1")
293+
set(RMW_UXRCE_MAX_SESSIONS "2")
294+
endif()
296295

297296
find_package(ament_lint_auto REQUIRED)
298297
ament_lint_auto_find_test_dependencies()
299298
add_subdirectory(test)
300299
endif()
301300

301+
# Create source files with the define
302+
configure_file(${PROJECT_SOURCE_DIR}/src/config.h.in
303+
${PROJECT_BINARY_DIR}/include/rmw_microxrcedds_c/config.h)
304+
302305
# Documentation
303306
if(BUILD_DOCUMENTATION)
304307
find_package(Doxygen)

rmw_microxrcedds_c/src/callbacks.c

Lines changed: 29 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -70,26 +70,31 @@ void on_topic(
7070
if ((custom_subscription->datareader_id.id == object_id.id) &&
7171
(custom_subscription->datareader_id.type == object_id.type))
7272
{
73-
rmw_uxrce_mempool_item_t * memory_node = get_memory(&static_buffer_memory);
73+
rmw_uxrce_mempool_item_t * memory_node = rmw_uxrce_get_static_input_buffer_for_entity(
74+
custom_subscription, custom_subscription->qos);
7475
if (!memory_node) {
7576
RMW_SET_ERROR_MSG("Not available static buffer memory node");
7677
return;
7778
}
7879

7980
rmw_uxrce_static_input_buffer_t * static_buffer =
8081
(rmw_uxrce_static_input_buffer_t *)memory_node->data;
81-
static_buffer->owner = (void *) custom_subscription;
82-
static_buffer->length = length;
8382

8483
if (!ucdr_deserialize_array_uint8_t(
8584
ub,
8685
static_buffer->buffer,
8786
length))
8887
{
8988
put_memory(&static_buffer_memory, memory_node);
89+
return;
9090
}
9191

92-
break;
92+
static_buffer->owner = (void *) custom_subscription;
93+
static_buffer->length = length;
94+
static_buffer->timestamp = rmw_uros_epoch_nanos();
95+
static_buffer->entity_type = RMW_UXRCE_ENTITY_TYPE_SUBSCRIPTION;
96+
97+
return;
9398
}
9499
subscription_item = subscription_item->next;
95100
}
@@ -114,27 +119,32 @@ void on_request(
114119
// Check if request is related to the service
115120
rmw_uxrce_service_t * custom_service = (rmw_uxrce_service_t *)service_item->data;
116121
if (custom_service->service_data_resquest == request_id) {
117-
rmw_uxrce_mempool_item_t * memory_node = get_memory(&static_buffer_memory);
122+
rmw_uxrce_mempool_item_t * memory_node = rmw_uxrce_get_static_input_buffer_for_entity(
123+
custom_service, custom_service->qos);
118124
if (!memory_node) {
119125
RMW_SET_ERROR_MSG("Not available static buffer memory node");
120126
return;
121127
}
122128

123129
rmw_uxrce_static_input_buffer_t * static_buffer =
124130
(rmw_uxrce_static_input_buffer_t *)memory_node->data;
125-
static_buffer->owner = (void *) custom_service;
126-
static_buffer->length = length;
127-
static_buffer->related.sample_id = *sample_id;
128131

129132
if (!ucdr_deserialize_array_uint8_t(
130133
ub,
131134
static_buffer->buffer,
132135
length))
133136
{
134137
put_memory(&static_buffer_memory, memory_node);
138+
return;
135139
}
136140

137-
break;
141+
static_buffer->owner = (void *) custom_service;
142+
static_buffer->length = length;
143+
static_buffer->related.sample_id = *sample_id;
144+
static_buffer->timestamp = rmw_uros_epoch_nanos();
145+
static_buffer->entity_type = RMW_UXRCE_ENTITY_TYPE_SERVICE;
146+
147+
return;
138148
}
139149
service_item = service_item->next;
140150
}
@@ -159,27 +169,32 @@ void on_reply(
159169
// Check if reply is related to the client
160170
rmw_uxrce_client_t * custom_client = (rmw_uxrce_client_t *)client_item->data;
161171
if (custom_client->client_data_request == request_id) {
162-
rmw_uxrce_mempool_item_t * memory_node = get_memory(&static_buffer_memory);
172+
rmw_uxrce_mempool_item_t * memory_node = rmw_uxrce_get_static_input_buffer_for_entity(
173+
custom_client, custom_client->qos);
163174
if (!memory_node) {
164175
RMW_SET_ERROR_MSG("Not available static buffer memory node");
165176
return;
166177
}
167178

168179
rmw_uxrce_static_input_buffer_t * static_buffer =
169180
(rmw_uxrce_static_input_buffer_t *)memory_node->data;
170-
static_buffer->owner = (void *) custom_client;
171-
static_buffer->length = length;
172-
static_buffer->related.reply_id = reply_id;
173181

174182
if (!ucdr_deserialize_array_uint8_t(
175183
ub,
176184
static_buffer->buffer,
177185
length))
178186
{
179187
put_memory(&static_buffer_memory, memory_node);
188+
return;
180189
}
181190

182-
break;
191+
static_buffer->owner = (void *) custom_client;
192+
static_buffer->length = length;
193+
static_buffer->related.reply_id = reply_id;
194+
static_buffer->timestamp = rmw_uros_epoch_nanos();
195+
static_buffer->entity_type = RMW_UXRCE_ENTITY_TYPE_CLIENT;
196+
197+
return;
183198
}
184199
client_item = client_item->next;
185200
}

rmw_microxrcedds_c/src/rmw_client.c

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ rmw_create_client(
7070
custom_client->rmw_handle = rmw_client;
7171
custom_client->owner_node = custom_node;
7272
custom_client->session_timeout = RMW_UXRCE_PUBLISH_RELIABLE_TIMEOUT;
73+
custom_client->qos = *qos_policies;
7374

7475
const rosidl_service_type_support_t * type_support_xrce = NULL;
7576
#ifdef ROSIDL_TYPESUPPORT_MICROXRCEDDS_C__IDENTIFIER_VALUE

rmw_microxrcedds_c/src/rmw_init.c

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -288,6 +288,14 @@ rmw_shutdown(
288288
*context = rmw_get_zero_initialized_context();
289289
}
290290

291+
rmw_uxrce_mempool_item_t * item = static_buffer_memory.allocateditems;
292+
293+
while (item != NULL) {
294+
rmw_uxrce_mempool_item_t * aux_next = item->next;
295+
put_memory(&static_buffer_memory, item);
296+
item = aux_next;
297+
}
298+
291299
return ret;
292300
}
293301

rmw_microxrcedds_c/src/rmw_publisher.c

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -102,8 +102,7 @@ rmw_create_publisher(
102102
custom_publisher->rmw_handle = rmw_publisher;
103103
custom_publisher->owner_node = custom_node;
104104
custom_publisher->session_timeout = RMW_UXRCE_PUBLISH_RELIABLE_TIMEOUT;
105-
106-
memcpy(&custom_publisher->qos, qos_policies, sizeof(rmw_qos_profile_t));
105+
custom_publisher->qos = *qos_policies;
107106

108107
custom_publisher->stream_id =
109108
(qos_policies->reliability == RMW_QOS_POLICY_RELIABILITY_BEST_EFFORT) ?
@@ -224,6 +223,7 @@ rmw_create_publisher(
224223
custom_publisher->topic->topic_id,
225224
reliability,
226225
history,
226+
qos_policies->depth,
227227
durability,
228228
UXR_REPLACE | UXR_REUSE);
229229
#endif /* ifdef RMW_UXRCE_USE_REFS */

rmw_microxrcedds_c/src/rmw_request.c

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,8 @@ rmw_take_request(
8181

8282
rmw_uxrce_service_t * custom_service = (rmw_uxrce_service_t *)service->data;
8383

84+
rmw_uxrce_clean_expired_static_input_buffer();
85+
8486
// Find first related item in static buffer memory pool
8587
rmw_uxrce_mempool_item_t * static_buffer_item =
8688
rmw_uxrce_find_static_input_buffer_by_owner((void *) custom_service);

rmw_microxrcedds_c/src/rmw_response.c

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,8 @@ rmw_take_response(
9595

9696
rmw_uxrce_client_t * custom_client = (rmw_uxrce_client_t *)client->data;
9797

98+
rmw_uxrce_clean_expired_static_input_buffer();
99+
98100
// Find first related item in static buffer memory pool
99101
rmw_uxrce_mempool_item_t * static_buffer_item =
100102
rmw_uxrce_find_static_input_buffer_by_owner((void *) custom_client);

rmw_microxrcedds_c/src/rmw_service.c

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@ rmw_create_service(
6969

7070
custom_service->owner_node = custom_node;
7171
custom_service->session_timeout = RMW_UXRCE_PUBLISH_RELIABLE_TIMEOUT;
72+
custom_service->qos = *qos_policies;
7273

7374
const rosidl_service_type_support_t * type_support_xrce = NULL;
7475
#ifdef ROSIDL_TYPESUPPORT_MICROXRCEDDS_C__IDENTIFIER_VALUE

rmw_microxrcedds_c/src/rmw_subscription.c

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -101,9 +101,8 @@ rmw_create_subscription(
101101

102102
rmw_uxrce_subscription_t * custom_subscription = (rmw_uxrce_subscription_t *)memory_node->data;
103103
custom_subscription->rmw_handle = rmw_subscription;
104-
105104
custom_subscription->owner_node = custom_node;
106-
memcpy(&custom_subscription->qos, qos_policies, sizeof(rmw_qos_profile_t));
105+
custom_subscription->qos = *qos_policies;
107106

108107
const rosidl_message_type_support_t * type_support_xrce = NULL;
109108
#ifdef ROSIDL_TYPESUPPORT_MICROXRCEDDS_C__IDENTIFIER_VALUE
@@ -210,6 +209,7 @@ rmw_create_subscription(
210209
custom_subscription->topic->topic_id,
211210
reliability,
212211
history,
212+
qos_policies->depth,
213213
durability,
214214
UXR_REPLACE | UXR_REUSE);
215215
#endif /* ifdef RMW_UXRCE_USE_XML */

rmw_microxrcedds_c/src/rmw_take.c

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,8 @@ rmw_take_with_info(
4949

5050
rmw_uxrce_subscription_t * custom_subscription = (rmw_uxrce_subscription_t *)subscription->data;
5151

52+
rmw_uxrce_clean_expired_static_input_buffer();
53+
5254
// Find first related item in static buffer memory pool
5355
rmw_uxrce_mempool_item_t * static_buffer_item = rmw_uxrce_find_static_input_buffer_by_owner(
5456
(void *) custom_subscription);

0 commit comments

Comments
 (0)