@@ -42,6 +42,7 @@ typedef struct app_data_t {
4242 int message_count ;
4343
4444 pn_proactor_t * proactor ;
45+ pn_message_t * message ;
4546 pn_rwbytes_t message_buffer ;
4647 int sent ;
4748 int acknowledged ;
@@ -58,40 +59,18 @@ static void check_condition(pn_event_t *e, pn_condition_t *cond) {
5859}
5960
6061/* Create a message with a map { "sequence" : number } encode it and return the encoded buffer. */
61- static pn_bytes_t encode_message (app_data_t * app ) {
62+ static void send_message (app_data_t * app , pn_link_t * sender ) {
6263 /* Construct a message with the map { "sequence": app.sent } */
63- pn_message_t * message = pn_message ();
64- pn_data_t * body = pn_message_body (message );
65- pn_message_set_id (message , (pn_atom_t ){.type = PN_ULONG , .u .as_ulong = app -> sent });
66- pn_data_put_map (body );
67- pn_data_enter (body );
68- pn_data_put_string (body , pn_bytes (sizeof ("sequence" )- 1 , "sequence" ));
69- pn_data_put_int (body , app -> sent ); /* The sequence number */
70- pn_data_exit (body );
71-
72- /* encode the message, expanding the encode buffer as needed */
73- if (app -> message_buffer .start == NULL ) {
74- static const size_t initial_size = 128 ;
75- app -> message_buffer = pn_rwbytes (initial_size , (char * )malloc (initial_size ));
76- }
77- /* app->message_buffer is the total buffer space available. */
78- /* mbuf wil point at just the portion used by the encoded message */
79- {
80- pn_rwbytes_t mbuf = pn_rwbytes (app -> message_buffer .size , app -> message_buffer .start );
81- int status = 0 ;
82- while ((status = pn_message_encode (message , mbuf .start , & mbuf .size )) == PN_OVERFLOW ) {
83- app -> message_buffer .size *= 2 ;
84- app -> message_buffer .start = (char * )realloc (app -> message_buffer .start , app -> message_buffer .size );
85- mbuf .size = app -> message_buffer .size ;
86- mbuf .start = app -> message_buffer .start ;
87- }
88- if (status != 0 ) {
89- fprintf (stderr , "error encoding message: %s\n" , pn_error_text (pn_message_error (message )));
64+ pn_bytes_t sequence = pn_bytes (sizeof ("sequence" )- 1 , "sequence" );
65+ pn_message_clear (app -> message );
66+ pn_message_set_id (app -> message , (pn_atom_t ){.type = PN_ULONG , .u .as_ulong = app -> sent });
67+ pn_amqp_map_t * props = pn_message_properties_build (NULL , sequence , (pn_atom_t ){.type = PN_INT , .u .as_int = app -> sent }, pn_bytes_null );
68+ pn_message_set_body_value (app -> message , (pn_amqp_value_t * )props );
69+ pn_amqp_map_free (props );
70+ if (pn_message_send (app -> message , sender , & app -> message_buffer ) < 0 ) {
71+ fprintf (stderr , "error sending message: %s\n" , pn_error_text (pn_message_error (app -> message )));
9072 exit (1 );
9173 }
92- pn_message_free (message );
93- return pn_bytes (mbuf .size , mbuf .start );
94- }
9574}
9675
9776/* Returns true to continue, false if finished */
@@ -140,11 +119,7 @@ static bool handle(app_data_t* app, pn_event_t* event) {
140119 ++ app -> sent ;
141120 /* Use sent counter as unique delivery tag. */
142121 pn_delivery (sender , pn_dtag ((const char * )& app -> sent , sizeof (app -> sent )));
143- {
144- pn_bytes_t msgbuf = encode_message (app );
145- pn_link_send (sender , msgbuf .start , msgbuf .size );
146- }
147- pn_link_advance (sender );
122+ send_message (app , sender );
148123 }
149124 break ;
150125 }
@@ -220,6 +195,7 @@ int main(int argc, char **argv) {
220195 app .port = (argc > 2 ) ? argv [2 ] : "amqp" ;
221196 app .amqp_address = (argc > 3 ) ? argv [3 ] : "examples" ;
222197 app .message_count = (argc > 4 ) ? atoi (argv [4 ]) : 10 ;
198+ app .message = pn_message ();
223199 app .user = (argc > 5 ) ? argv [5 ] : 0 ;
224200 app .pass = (argc > 6 ) ? argv [6 ] : 0 ;
225201
@@ -250,5 +226,6 @@ int main(int argc, char **argv) {
250226
251227 pn_proactor_free (app .proactor );
252228 free (app .message_buffer .start );
229+ pn_message_free (app .message );
253230 return exit_code ;
254231}
0 commit comments