Logo AND Algorithmique Numérique Distribuée

Public GIT Repository
also awake the listener after opening a file socket
[simgrid.git] / src / gras / Transport / transport_plugin_file.c
index 7e7eae0..79e8cfc 100644 (file)
@@ -8,26 +8,29 @@
  * under the terms of the license (GNU LGPL) which comes with this package. */
 
 #include "portable.h"
-#include "transport_private.h"
+#include "gras/Transport/transport_private.h"
 #include "xbt/ex.h"
+#include "gras/Msg/msg_interface.h" /* gras_msg_listener_awake */
 
-XBT_LOG_NEW_DEFAULT_SUBCATEGORY(trp_file,transport,
+XBT_LOG_NEW_DEFAULT_SUBCATEGORY(gras_trp_file,gras_trp,
        "Pseudo-transport to write to/read from a file");
 
 /***
- *** Prototypes 
+ *** Prototypes
  ***/
 void gras_trp_file_close(gras_socket_t sd);
-  
+
+void gras_trp_file_chunk_send_raw(gras_socket_t sd,
+                                 const char *data,
+                                 unsigned long int size);
 void gras_trp_file_chunk_send(gras_socket_t sd,
                              const char *data,
-                             unsigned long int size);
-
-void gras_trp_file_chunk_recv(gras_socket_t sd,
-                             char *data,
                              unsigned long int size,
-                             unsigned long int bufsize);
+                             int stable_ignored);
 
+int gras_trp_file_chunk_recv(gras_socket_t sd,
+                            char *data,
+                            unsigned long int size);
 
 /***
  *** Specific plugin part
@@ -54,8 +57,12 @@ gras_trp_file_setup(gras_trp_plugin_t plug) {
   FD_ZERO(&(file->incoming_socks));
 
   plug->socket_close = gras_trp_file_close;
-  plug->chunk_send   = gras_trp_file_chunk_send;
-  plug->chunk_recv   = gras_trp_file_chunk_recv;
+
+  plug->raw_send = gras_trp_file_chunk_send_raw;
+  plug->send = gras_trp_file_chunk_send;
+
+  plug->raw_recv = plug->recv = gras_trp_file_chunk_recv;
+
   plug->data         = (void*)file;
 }
 
@@ -78,8 +85,8 @@ gras_socket_client_from_file(const char*path) {
   res->plugin=gras_trp_plugin_get_by_name("file");
 
   if (strcmp("-", path)) {
-    res->sd = open(path, O_WRONLY|O_CREAT | O_BINARY, S_IRUSR|S_IWUSR|S_IRGRP );
-    
+    res->sd = open(path, O_TRUNC|O_WRONLY|O_CREAT | O_BINARY, S_IRUSR|S_IWUSR|S_IRGRP );
+
     if ( res->sd < 0) {
       THROW2(system_error,0,
             "Cannot create a client socket from file %s: %s",
@@ -92,10 +99,12 @@ gras_socket_client_from_file(const char*path) {
   DEBUG5("sock_client_from_file(%s): sd=%d in=%c out=%c accept=%c",
         path,
         res->sd,
-        res->incoming?'y':'n', 
+        res->incoming?'y':'n',
         res->outgoing?'y':'n',
         res->accepting?'y':'n');
 
+  xbt_dynar_push(((gras_trp_procdata_t)
+                 gras_libdata_by_id(gras_trp_libdata_id))->sockets,&res);
   return res;
 }
 
@@ -131,16 +140,19 @@ gras_socket_t gras_socket_server_from_file(const char*path) {
 
   DEBUG4("sd=%d in=%c out=%c accept=%c",
         res->sd,
-        res->incoming?'y':'n', 
+        res->incoming?'y':'n',
         res->outgoing?'y':'n',
         res->accepting?'y':'n');
 
+  xbt_dynar_push(((gras_trp_procdata_t)
+                 gras_libdata_by_id(gras_trp_libdata_id))->sockets,&res);
+  gras_msg_listener_awake();
   return res;
 }
 
 void gras_trp_file_close(gras_socket_t sock){
   gras_trp_file_plug_data_t *data;
-  
+
   if (!sock) return; /* close only once */
   data=sock->plugin->data;
 
@@ -156,7 +168,7 @@ void gras_trp_file_close(gras_socket_t sock){
 
     /* close the socket */
     if(close(sock->sd) < 0) {
-      WARN2("error while closing file %d: %s", 
+      WARN2("error while closing file %d: %s",
               sock->sd, strerror(errno));
     }
   }
@@ -170,23 +182,30 @@ void gras_trp_file_close(gras_socket_t sock){
 void
 gras_trp_file_chunk_send(gras_socket_t sock,
                         const char *data,
-                        unsigned long int size) {
-  
+                        unsigned long int size,
+                        int stable_ignored) {
+  gras_trp_file_chunk_send_raw(sock,data,size);
+}
+void
+gras_trp_file_chunk_send_raw(gras_socket_t sock,
+                            const char *data,
+                            unsigned long int size) {
+
   xbt_assert0(sock->outgoing, "Cannot write on client file socket");
   xbt_assert0(size >= 0, "Cannot send a negative amount of data");
 
   while (size) {
     int status = 0;
-    
+
     DEBUG3("write(%d, %p, %ld);", sock->sd, data, (long int)size);
     status = write(sock->sd, data, (long int)size);
-    
+
     if (status == -1) {
       THROW4(system_error,0,"write(%d,%p,%d) failed: %s",
             sock->sd, data, (int)size,
             strerror(errno));
     }
-    
+
     if (status) {
       size  -= status;
       data  += status;
@@ -200,36 +219,43 @@ gras_trp_file_chunk_send(gras_socket_t sock,
  *
  * Receive data on a file pseudo-socket.
  */
-void
+int
 gras_trp_file_chunk_recv(gras_socket_t sock,
                         char *data,
-                        unsigned long int size,
-                        unsigned long int bufsize) {
+                        unsigned long int size) {
+
+  int got = 0;
 
   xbt_assert0(sock, "Cannot recv on an NULL socket");
   xbt_assert0(sock->incoming, "Cannot recv on client file socket");
   xbt_assert0(size >= 0, "Cannot receive a negative amount of data");
-  xbt_assert0(bufsize>=size,"Not enough buffer size to receive that much data");
+
+  if (sock->recvd) {
+     data[0] = sock->recvd_val;
+     sock->recvd = 0;
+     got++;
+     size--;
+  }
 
   while (size) {
     int status = 0;
-    
-    status = read(sock->sd, data, (long int)bufsize);
-    DEBUG3("read(%d, %p, %ld);", sock->sd, data, size);
-    
-    if (status == -1) {
+
+    status = read(sock->sd, data+got, (long int)size);
+    DEBUG3("read(%d, %p, %ld);", sock->sd, data+got, size);
+
+    if (status < 0) {
       THROW4(system_error,0,"read(%d,%p,%d) failed: %s",
-            sock->sd, data, (int)size,
+            sock->sd, data+got, (int)size,
             strerror(errno));
     }
-    
+
     if (status) {
       size    -= status;
-      bufsize -= status;
-      data    += status;
+      got    += status;
     } else {
-      THROW0(system_error,0,"file descriptor closed");
+       THROW1(system_error,errno,"file descriptor closed after %d bytes",got);
     }
   }
+  return got;
 }