- asprintf(&send_msg, "get_subbuffer %s", buf->name);
- result = ustcomm_send_request(&buf->conn, send_msg, &received_msg);
- if(test_sigpipe()) {
- WARN("process %d destroyed before we could connect to it", buf->pid);
- retval = GET_SUBBUF_DONE;
- goto end;
- }
- else if(result < 0) {
- ERR("get_subbuffer: ustcomm_send_request failed");
- retval = -1;
- goto end;
- }
-
- result = sscanf(received_msg, "%as %ld", &rep_code, &buf->consumed_old);
- if(result != 2 && result != 1) {
- ERR("unable to parse response to get_subbuffer");
- retval = -1;
- goto end_rep;
- }
-
- DBG("received msg is %s", received_msg);
-
- if(!strcmp(rep_code, "OK")) {
- DBG("got subbuffer %s", buf->name);
- retval = GET_SUBBUF_OK;
- }
- else if(nth_token_is(received_msg, "END", 0) == 1) {
- retval = GET_SUBBUF_DONE;
- goto end_rep;
- }
- else {
- DBG("error getting subbuffer %s", buf->name);
- retval = -1;
- }
-
- /* FIMXE: free correctly the stuff */
-end_rep:
- if(rep_code)
- free(rep_code);
-end:
- if(send_msg)
- free(send_msg);
- if(received_msg)
- free(received_msg);
-
- return retval;
-}
-
-int put_subbuffer(struct buffer_info *buf)
-{
- char *send_msg=NULL;
- char *received_msg=NULL;
- char *rep_code=NULL;
- int retval;
- int result;
-
- asprintf(&send_msg, "put_subbuffer %s %ld", buf->name, buf->consumed_old);
- result = ustcomm_send_request(&buf->conn, send_msg, &received_msg);
- if(result < 0 && errno == ECONNRESET) {
- retval = PUT_SUBBUF_DIED;
- goto end;
- }
- if(result < 0) {
- ERR("put_subbuffer: send_message failed");
- retval = -1;
- goto end;
- }
-
- result = sscanf(received_msg, "%as", &rep_code);
- if(result != 1) {
- ERR("unable to parse response to put_subbuffer");
- retval = -1;
- goto end_rep;
- }
-
- if(!strcmp(rep_code, "OK")) {
- DBG("subbuffer put %s", buf->name);
- retval = PUT_SUBBUF_OK;
- }
- else {
- DBG("put_subbuffer: received error, we were pushed");
- retval = PUT_SUBBUF_PUSHED;
- goto end_rep;
- }
-
-end_rep:
- if(rep_code)
- free(rep_code);
-
-end:
- if(send_msg)
- free(send_msg);
- if(received_msg)
- free(received_msg);
-
- return retval;
-}
-
-/* This write is patient because it restarts if it was incomplete.
- */
-
-ssize_t patient_write(int fd, const void *buf, size_t count)
-{
- const char *bufc = (const char *) buf;
- int result;
-
- for(;;) {
- result = write(fd, bufc, count);
- if(result <= 0) {
- return result;
- }
- count -= result;
- bufc += result;
-
- if(count == 0) {
- break;
- }
- }
-
- return bufc-(const char *)buf;
-}
-
-void decrement_active_buffers(void *arg)
-{
- pthread_mutex_lock(&active_buffers_mutex);
- active_buffers--;
- pthread_mutex_unlock(&active_buffers_mutex);
-}
-
-void *consumer_thread(void *arg)
-{
- struct buffer_info *buf = (struct buffer_info *) arg;
- int result;
-
- pthread_cleanup_push(decrement_active_buffers, NULL);
-
- for(;;) {
- /* get the subbuffer */
- result = get_subbuffer(buf);
- if(result == -1) {
- ERR("error getting subbuffer");
- continue;
- }
- else if(result == GET_SUBBUF_DONE) {
- /* this is done */
- break;
- }
- else if(result == GET_SUBBUF_DIED) {
- finish_consuming_dead_subbuffer(buf);
- break;
- }
-
- /* write data to file */
- result = patient_write(buf->file_fd, buf->mem + (buf->consumed_old & (buf->n_subbufs * buf->subbuf_size-1)), buf->subbuf_size);
- if(result == -1) {
- PERROR("write");
- /* FIXME: maybe drop this trace */
- }
-
- /* put the subbuffer */
- result = put_subbuffer(buf);
- if(result == -1) {
- ERR("unknown error putting subbuffer (channel=%s)", buf->name);
- break;
- }
- else if(result == PUT_SUBBUF_PUSHED) {
- ERR("Buffer overflow (channel=%s), reader pushed. This channel will not be usable passed this point.", buf->name);
- break;
- }
- else if(result == PUT_SUBBUF_DIED) {
- WARN("application died while putting subbuffer");
- /* FIXME: probably need to skip the first subbuffer in finish_consuming_dead_subbuffer */
- finish_consuming_dead_subbuffer(buf);
- break;
- }
- else if(result == PUT_SUBBUF_OK) {
- }