void amok_bw_init(void) {
gras_datadesc_type_t bw_request_desc, bw_res_desc, sat_request_desc;
- if (_amok_bw_initialized)
- return;
-
amok_base_init();
-
- /* Build the datatype descriptions */
- bw_request_desc = gras_datadesc_struct("s_bw_request_t");
- gras_datadesc_struct_append(bw_request_desc,"host",gras_datadesc_by_name("xbt_host_t"));
- gras_datadesc_struct_append(bw_request_desc,"buf_size",gras_datadesc_by_name("unsigned int"));
- gras_datadesc_struct_append(bw_request_desc,"exp_size",gras_datadesc_by_name("unsigned int"));
- gras_datadesc_struct_append(bw_request_desc,"msg_size",gras_datadesc_by_name("unsigned int"));
- gras_datadesc_struct_close(bw_request_desc);
- bw_request_desc = gras_datadesc_ref("bw_request_t",bw_request_desc);
-
- bw_res_desc = gras_datadesc_struct("s_bw_res_t");
- gras_datadesc_struct_append(bw_res_desc,"err",gras_datadesc_by_name("s_amok_remoterr_t"));
- gras_datadesc_struct_append(bw_res_desc,"timestamp",gras_datadesc_by_name("unsigned int"));
- gras_datadesc_struct_append(bw_res_desc,"seconds",gras_datadesc_by_name("double"));
- gras_datadesc_struct_append(bw_res_desc,"bw",gras_datadesc_by_name("double"));
- gras_datadesc_struct_close(bw_res_desc);
- bw_res_desc = gras_datadesc_ref("bw_res_t",bw_res_desc);
-
- sat_request_desc = gras_datadesc_struct("s_sat_request_desc_t");
- gras_datadesc_struct_append(sat_request_desc,"host",gras_datadesc_by_name("xbt_host_t"));
- gras_datadesc_struct_append(sat_request_desc,"msg_size",gras_datadesc_by_name("unsigned int"));
- gras_datadesc_struct_append(sat_request_desc,"timeout",gras_datadesc_by_name("unsigned int"));
- gras_datadesc_struct_close(sat_request_desc);
- sat_request_desc = gras_datadesc_ref("sat_request_t",sat_request_desc);
-
- /* Register the bandwidth messages */
- gras_msgtype_declare("BW request", bw_request_desc);
- gras_msgtype_declare("BW result", bw_res_desc);
- gras_msgtype_declare("BW handshake", bw_request_desc);
- gras_msgtype_declare("BW handshake ACK", bw_request_desc);
-
- /* Register the saturation messages */
- gras_msgtype_declare("SAT start", sat_request_desc);
- gras_msgtype_declare("SAT started", gras_datadesc_by_name("amok_remoterr_t"));
- gras_msgtype_declare("SAT begin", sat_request_desc);
- gras_msgtype_declare("SAT begun", gras_datadesc_by_name("amok_remoterr_t"));
- gras_msgtype_declare("SAT end", NULL);
- gras_msgtype_declare("SAT ended", gras_datadesc_by_name("amok_remoterr_t"));
- gras_msgtype_declare("SAT stop", NULL);
- gras_msgtype_declare("SAT stopped", gras_datadesc_by_name("amok_remoterr_t"));
+ if (! _amok_bw_initialized) {
+
+ /* Build the datatype descriptions */
+ bw_request_desc = gras_datadesc_struct("s_bw_request_t");
+ gras_datadesc_struct_append(bw_request_desc,"host",gras_datadesc_by_name("xbt_host_t"));
+ gras_datadesc_struct_append(bw_request_desc,"buf_size",gras_datadesc_by_name("unsigned int"));
+ gras_datadesc_struct_append(bw_request_desc,"exp_size",gras_datadesc_by_name("unsigned int"));
+ gras_datadesc_struct_append(bw_request_desc,"msg_size",gras_datadesc_by_name("unsigned int"));
+ gras_datadesc_struct_close(bw_request_desc);
+ bw_request_desc = gras_datadesc_ref("bw_request_t",bw_request_desc);
+
+ bw_res_desc = gras_datadesc_struct("s_bw_res_t");
+ gras_datadesc_struct_append(bw_res_desc,"err",gras_datadesc_by_name("s_amok_remoterr_t"));
+ gras_datadesc_struct_append(bw_res_desc,"timestamp",gras_datadesc_by_name("unsigned int"));
+ gras_datadesc_struct_append(bw_res_desc,"seconds",gras_datadesc_by_name("double"));
+ gras_datadesc_struct_append(bw_res_desc,"bw",gras_datadesc_by_name("double"));
+ gras_datadesc_struct_close(bw_res_desc);
+ bw_res_desc = gras_datadesc_ref("bw_res_t",bw_res_desc);
+
+ sat_request_desc = gras_datadesc_struct("s_sat_request_desc_t");
+ gras_datadesc_struct_append(sat_request_desc,"host",gras_datadesc_by_name("xbt_host_t"));
+ gras_datadesc_struct_append(sat_request_desc,"msg_size",gras_datadesc_by_name("unsigned int"));
+ gras_datadesc_struct_append(sat_request_desc,"timeout",gras_datadesc_by_name("unsigned int"));
+ gras_datadesc_struct_close(sat_request_desc);
+ sat_request_desc = gras_datadesc_ref("sat_request_t",sat_request_desc);
+
+ /* Register the bandwidth messages */
+ gras_msgtype_declare("BW request", bw_request_desc);
+ gras_msgtype_declare("BW result", bw_res_desc);
+ gras_msgtype_declare("BW handshake", bw_request_desc);
+ gras_msgtype_declare("BW handshake ACK", bw_request_desc);
+
+ /* Register the saturation messages */
+ gras_msgtype_declare("SAT start", sat_request_desc);
+ gras_msgtype_declare("SAT started", gras_datadesc_by_name("amok_remoterr_t"));
+ gras_msgtype_declare("SAT begin", sat_request_desc);
+ gras_msgtype_declare("SAT begun", gras_datadesc_by_name("amok_remoterr_t"));
+ gras_msgtype_declare("SAT end", NULL);
+ gras_msgtype_declare("SAT ended", gras_datadesc_by_name("amok_remoterr_t"));
+ gras_msgtype_declare("SAT stop", NULL);
+ gras_msgtype_declare("SAT stopped", gras_datadesc_by_name("amok_remoterr_t"));
+ }
+
/* Register the callbacks */
gras_cb_register(gras_msgtype_by_name("BW request"),
&amok_bw_cb_bw_request);
xbt_error_t amok_bw_test(gras_socket_t peer,
unsigned int buf_size,unsigned int exp_size,unsigned int msg_size,
/*OUT*/ double *sec, double *bw) {
- gras_socket_t rawIn,rawOut; /* raw sockets for the experiments */
+ gras_socket_t measIn,measOut; /* Measurement sockets for the experiments */
gras_socket_t sock_dummy; /* ignored arg to msg_wait */
int port;
xbt_error_t errcode;
for (port = 5000, errcode = system_error;
errcode == system_error;
- errcode = gras_socket_server_ext(++port,buf_size,1,&rawIn));
+ errcode = gras_socket_server_ext(++port,buf_size,1,&measIn));
if (errcode != no_error) {
- ERROR1("Error %s encountered while opening a raw socket",
+ ERROR1("Error %s encountered while opening a measurement socket",
xbt_error_name(errcode));
return errcode;
}
request->exp_size=exp_size;
request->msg_size=msg_size;
request->host.name = NULL;
- request->host.port = gras_socket_my_port(rawIn);
- INFO1("Send an handshake to get the dude connect to port %d on me", request->host.port);
+ request->host.port = gras_socket_my_port(measIn);
+ INFO3("Handshaking with %s:%d to connect it back on my %d",
+ gras_socket_peer_name(peer),gras_socket_peer_port(peer), request->host.port);
if ((errcode=gras_msg_send(peer,gras_msgtype_by_name("BW handshake"),&request))) {
ERROR1("Error %s encountered while sending the BW request.", xbt_error_name(errcode));
/* FIXME: What if there is a remote error? */
- if((errcode=gras_socket_client_ext(gras_socket_peer_name(peer),request_ack->host.port, buf_size,1,&rawOut))) {
- ERROR3("Error %s encountered while opening the raw socket to %s:%d for BW test\n",
+ if((errcode=gras_socket_client_ext(gras_socket_peer_name(peer),request_ack->host.port, buf_size,1,&measOut))) {
+ ERROR3("Error %s encountered while opening the measurement socket to %s:%d for BW test\n",
xbt_error_name(errcode),gras_socket_peer_name(peer),request_ack->host.port);
return errcode;
}
INFO0("Got ACK");
*sec=gras_os_time();
- if ((errcode=gras_socket_raw_send(rawOut,120,exp_size,msg_size)) ||
- (errcode=gras_socket_raw_recv(rawIn,120,1,1))) {
+ if ((errcode=gras_socket_meas_send(measOut,120,exp_size,msg_size)) ||
+ (errcode=gras_socket_meas_recv(measIn,120,1,1))) {
ERROR1("Error %s encountered while sending the BW experiment.",
xbt_error_name(errcode));
- gras_socket_close(rawOut);
- gras_socket_close(rawIn);
+ gras_socket_close(measOut);
+ gras_socket_close(measIn);
return errcode;
}
*sec = gras_os_time() - *sec;
*bw = ((double)exp_size /* 8.0*/) / *sec / (1024.0 *1024.0);
INFO0("DOOONE");
- gras_socket_close(rawIn);
- gras_socket_close(rawOut);
+ gras_socket_close(measIn);
+ gras_socket_close(measOut);
return no_error;
}
/* Callback to the "BW handshake" message:
- opens a server raw socket,
+ opens a server measurement socket,
indicate its port in an "BW handshaked" message,
- receive the corresponding data on the raw socket,
- close the raw socket
+ receive the corresponding data on the measurement socket,
+ close the measurment socket
*/
int amok_bw_cb_bw_handshake(gras_socket_t expeditor,
void *payload) {
- gras_socket_t rawIn,rawOut;
+ gras_socket_t measIn,measOut;
bw_request_t request=*(bw_request_t*)payload;
bw_request_t answer;
xbt_error_t errcode;
int port;
- INFO2("Got an handshake to connect to %s:%d",
- request->host.name,request->host.port);
+ INFO2("Handshaked to connect to %s:%d",
+ gras_socket_peer_name(expeditor),request->host.port);
answer = xbt_new0(s_bw_request_t,1);
- for (port = 5000, errcode = system_error;
+ for (port = 6000, errcode = system_error;
errcode == system_error;
- errcode = gras_socket_server_ext(++port,request->buf_size,1,&rawIn));
+ errcode = gras_socket_server_ext(++port,request->buf_size,1,&measIn));
if (errcode != no_error) {
- ERROR1("Error %s encountered while opening a raw socket", xbt_error_name(errcode));
+ ERROR1("Error %s encountered while opening a measurement server socket", xbt_error_name(errcode));
/* FIXME: tell error to remote */
return 1;
}
if ((errcode=gras_socket_client_ext(gras_socket_peer_name(expeditor),request->host.port,
- request->buf_size,1,&rawOut))) {
- ERROR3("Error '%s' encountered while opening a raw socket to %s:%d",
+ request->buf_size,1,&measOut))) {
+ ERROR3("Error '%s' encountered while opening a measurement socket back to %s:%d",
xbt_error_name(errcode),gras_socket_peer_name(expeditor),request->host.port);
/* FIXME: tell error to remote */
return 1;
answer->buf_size=request->buf_size;
answer->exp_size=request->exp_size;
answer->msg_size=request->msg_size;
- answer->host.port=gras_socket_my_port(rawIn);
+ answer->host.port=gras_socket_my_port(measIn);
if ((errcode=gras_msg_send(expeditor,gras_msgtype_by_name("BW handshake ACK"),&answer))) {
ERROR1("Error %s encountered while sending the answer.",
xbt_error_name(errcode));
- gras_socket_close(rawIn);
- gras_socket_close(rawOut);
+ gras_socket_close(measIn);
+ gras_socket_close(measOut);
/* FIXME: tell error to remote */
return 1;
}
INFO4("BW handshake answered. buf_size=%d exp_size=%d msg_size=%d port=%d",
answer->buf_size,answer->exp_size,answer->msg_size,answer->host.port);
- if ((errcode=gras_socket_raw_recv(rawIn, 120,request->exp_size,request->msg_size)) ||
- (errcode=gras_socket_raw_send(rawOut,120,1,1))) {
+ if ((errcode=gras_socket_meas_recv(measIn, 120,request->exp_size,request->msg_size)) ||
+ (errcode=gras_socket_meas_send(measOut,120,1,1))) {
ERROR1("Error %s encountered while receiving the experiment.",
xbt_error_name(errcode));
- gras_socket_close(rawIn);
- gras_socket_close(rawOut);
+ gras_socket_close(measIn);
+ gras_socket_close(measOut);
/* FIXME: tell error to remote ? */
return 1;
}
- gras_socket_close(rawIn);
- gras_socket_close(rawOut);
+ gras_socket_close(measIn);
+ gras_socket_close(measOut);
return 1;
}
int amok_bw_cb_bw_request(gras_socket_t expeditor,
- void *payload) {return 1;}
+ void *payload) {
+ CRITICAL0("Not implemented");
+ return 1;
+}
int amok_bw_cb_sat_start(gras_socket_t expeditor,
- void *payload) {return 1;}
+ void *payload) {
+ CRITICAL0("Not implemented");
+ return 1;
+}
int amok_bw_cb_sat_begin(gras_socket_t expeditor,
- void *payload) {return 1;}
+ void *payload) {
+ CRITICAL0("Not implemented");
+ return 1;
+}
#if 0
/* function to request a BW test between two external hosts */