3 /* buf trp (transport) - buffered transport using the TCP one */
5 /* Copyright (c) 2004 Martin Quinson. All rights reserved. */
7 /* This program is free software; you can redistribute it and/or modify it
8 * under the terms of the license (GNU LGPL) which comes with this package. */
10 #include <netinet/in.h> /* htonl/ntohl */
12 #include <string.h> /* memset */
15 #include "transport_private.h"
17 XBT_LOG_NEW_DEFAULT_SUBCATEGORY(trp_buf,transport,
18 "Generic buffered transport (works on top of TCP or SG)");
23 xbt_error_t gras_trp_buf_socket_client(gras_trp_plugin_t *self,
25 xbt_error_t gras_trp_buf_socket_server(gras_trp_plugin_t *self,
27 xbt_error_t gras_trp_buf_socket_accept(gras_socket_t sock,
30 void gras_trp_buf_socket_close(gras_socket_t sd);
32 xbt_error_t gras_trp_buf_chunk_send(gras_socket_t sd,
36 xbt_error_t gras_trp_buf_chunk_recv(gras_socket_t sd,
39 xbt_error_t gras_trp_buf_flush(gras_socket_t sock);
43 *** Specific plugin part
47 gras_trp_plugin_t *super;
48 } gras_trp_buf_plug_data_t;
51 *** Specific socket part
57 int pos; /* for receive; not exchanged over the net */
60 struct gras_trp_bufdata_{
66 void gras_trp_buf_init_sock(gras_socket_t sock) {
67 gras_trp_bufdata_t *data=xbt_new(gras_trp_bufdata_t,1);
70 data->buffsize = 100 * 1024 ; /* 100k */
73 data->in.data = xbt_malloc(data->buffsize);
74 data->in.pos = 0; /* useless, indeed, since size==pos */
77 data->out.data = xbt_malloc(data->buffsize);
88 gras_trp_buf_setup(gras_trp_plugin_t *plug) {
90 gras_trp_buf_plug_data_t *data =xbt_new(gras_trp_buf_plug_data_t,1);
93 TRY(gras_trp_plugin_get_by_name(gras_if_RL() ? "tcp" : "sg",
95 DEBUG1("Derivate a buffer plugin from %s",gras_if_RL() ? "tcp" : "sg");
97 plug->socket_client = gras_trp_buf_socket_client;
98 plug->socket_server = gras_trp_buf_socket_server;
99 plug->socket_accept = gras_trp_buf_socket_accept;
100 plug->socket_close = gras_trp_buf_socket_close;
102 plug->chunk_send = gras_trp_buf_chunk_send;
103 plug->chunk_recv = gras_trp_buf_chunk_recv;
105 plug->flush = gras_trp_buf_flush;
107 plug->data = (void*)data;
113 xbt_error_t gras_trp_buf_socket_client(gras_trp_plugin_t *self,
114 /* OUT */ gras_socket_t sock){
116 gras_trp_plugin_t *super=((gras_trp_buf_plug_data_t*)self->data)->super;
119 TRY(super->socket_client(super,sock));
121 gras_trp_buf_init_sock(sock);
127 * gras_trp_buf_socket_server:
129 * Open a socket used to receive messages.
131 xbt_error_t gras_trp_buf_socket_server(gras_trp_plugin_t *self,
132 /* OUT */ gras_socket_t sock){
134 gras_trp_plugin_t *super=((gras_trp_buf_plug_data_t*)self->data)->super;
137 TRY(super->socket_server(super,sock));
139 gras_trp_buf_init_sock(sock);
144 gras_trp_buf_socket_accept(gras_socket_t sock,
145 gras_socket_t *dst) {
147 gras_trp_plugin_t *super=((gras_trp_buf_plug_data_t*)sock->plugin->data)->super;
150 TRY(super->socket_accept(sock,dst));
151 (*dst)->plugin = sock->plugin;
152 gras_trp_buf_init_sock(*dst);
156 void gras_trp_buf_socket_close(gras_socket_t sock){
157 gras_trp_plugin_t *super=((gras_trp_buf_plug_data_t*)sock->plugin->data)->super;
158 gras_trp_bufdata_t *data=sock->bufdata;
161 if (data->in.size || data->out.size)
162 gras_trp_buf_flush(sock);
164 xbt_free(data->in.data);
166 xbt_free(data->out.data);
169 super->socket_close(sock);
173 * gras_trp_buf_chunk_send:
175 * Send data on a TCP socket
178 gras_trp_buf_chunk_send(gras_socket_t sock,
183 gras_trp_bufdata_t *data=(gras_trp_bufdata_t*)sock->bufdata;
187 /* Let underneath plugin check for direction, we work even in duplex */
188 xbt_assert0(size >= 0, "Cannot send a negative amount of data");
190 while (chunk_pos < size) {
191 /* size of the chunck to receive in that shot */
192 long int thissize = min(size-chunk_pos,data->buffsize - data->out.size);
193 DEBUG5("Set the chars %d..%ld into the buffer (size=%ld, ctn='%.*s')",
195 ((int)data->out.size) + thissize -1,
196 size, chunk_pos, chunk);
198 memcpy(data->out.data + data->out.size, chunk + chunk_pos, thissize);
200 data->out.size += thissize;
201 chunk_pos += thissize;
202 DEBUG5("New pos = %d; Still to send = %ld of %ld; ctn sofar='%.*s'",
203 data->out.size,size-chunk_pos,size,(int)chunk_pos,chunk);
205 if (data->out.size == data->buffsize) /* out of space. Flush it */
206 TRY(gras_trp_buf_flush(sock));
214 * gras_trp_buf_chunk_recv:
216 * Receive data on a TCP socket.
219 gras_trp_buf_chunk_recv(gras_socket_t sock,
224 gras_trp_plugin_t *super=((gras_trp_buf_plug_data_t*)sock->plugin->data)->super;
225 gras_trp_bufdata_t *data=sock->bufdata;
226 long int chunck_pos = 0;
228 /* Let underneath plugin check for direction, we work even in duplex */
229 xbt_assert0(sock, "Cannot recv on an NULL socket");
230 xbt_assert0(size >= 0, "Cannot receive a negative amount of data");
234 while (chunck_pos < size) {
235 /* size of the chunck to receive in that shot */
238 if (data->in.size == data->in.pos) { /* out of data. Get more */
240 DEBUG0("Recv the size");
241 TRY(super->chunk_recv(sock,(char*)&nextsize, 4));
242 data->in.size = ntohl(nextsize);
244 VERB1("Recv the chunk (size=%d)",data->in.size);
245 TRY(super->chunk_recv(sock, data->in.data, data->in.size));
249 thissize = min(size-chunck_pos , data->in.size - data->in.pos);
250 DEBUG2("Get the chars %d..%ld out of the buffer",
252 data->in.pos + thissize - 1);
253 memcpy(chunk+chunck_pos, data->in.data + data->in.pos, thissize);
255 data->in.pos += thissize;
256 chunck_pos += thissize;
257 DEBUG5("New pos = %d; Still to receive = %ld of %ld. Ctn so far='%.*s'",
258 data->in.pos,size - chunck_pos,size,(int)chunck_pos,chunk);
266 * gras_trp_buf_flush:
268 * Make sure the data is sent
271 gras_trp_buf_flush(gras_socket_t sock) {
274 gras_trp_plugin_t *super=((gras_trp_buf_plug_data_t*)sock->plugin->data)->super;
275 gras_trp_bufdata_t *data=sock->bufdata;
278 size = htonl(data->out.size);
279 DEBUG1("Send the size (=%d)",data->out.size);
280 TRY(super->chunk_send(sock,(char*) &size, 4));
282 DEBUG1("Send the chunk (size=%d)",data->out.size);
283 TRY(super->chunk_send(sock, data->out.data, data->out.size));
284 VERB1("Chunk sent (size=%d)",data->out.size);