| 7732 | return job; |
| 7733 | } |
| 7734 | |
| 7735 | static void *dist_worker_prefetch_eval_main(void *arg) { |
| 7736 | ds4_dist_worker_job_queue *q = arg; |
| 7737 | for (;;) { |
| 7738 | ds4_dist_worker_job *job = dist_worker_job_queue_pop(q); |
| 7739 | if (!job) break; |
| 7740 | int rc = dist_worker_process_work_payload(q->state, |
| 7741 | q->upstream, |
| 7742 | job->payload, |
| 7743 | job->bytes); |
| 7744 | dist_worker_job_free(job); |
| 7745 | if (rc <= 0) { |
| 7746 | pthread_mutex_lock(&q->mu); |
| 7747 | q->rc = rc == 0 ? 0 : 1; |
| 7748 | pthread_mutex_unlock(&q->mu); |
| 7749 | dist_worker_job_queue_cancel(q); |
| 7750 | shutdown(q->upstream->fd, SHUT_RDWR); |
| 7751 | break; |
| 7752 | } |
| 7753 | } |
| 7754 | return NULL; |
| 7755 | } |
| 7756 | |
| 7757 | static int dist_worker_read_loop_prefetch(ds4_dist_worker_state *state, int fd) { |
| 7758 | ds4_dist_worker_upstream upstream; |
| 7759 | dist_worker_upstream_init(&upstream, state, fd); |
| 7760 | |
| 7761 | ds4_dist_worker_job_queue queue; |
| 7762 | dist_worker_job_queue_init(&queue, state, &upstream); |
| 7763 | |
| 7764 | pthread_t eval_tid; |
| 7765 | if (pthread_create(&eval_tid, NULL, dist_worker_prefetch_eval_main, &queue) != 0) { |
| 7766 | dist_worker_job_queue_destroy(&queue); |
| 7767 | dist_worker_upstream_destroy(&upstream); |
| 7768 | return 1; |
| 7769 | } |
| 7770 | |
| 7771 | int loop_rc = 0; |
| 7772 | fprintf(stderr, |
| 7773 | "ds4: distributed worker: receive prefetch depth %u enabled\n", |
| 7774 | queue.depth); |
| 7775 | |
| 7776 | for (;;) { |
| 7777 | uint32_t type = 0, bytes = 0; |
| 7778 | char err[256]; |
| 7779 | int rc = dist_read_frame_header(fd, &type, &bytes, err, sizeof(err)); |
| 7780 | if (rc == 0) break; |
| 7781 | if (rc < 0) { |
| 7782 | fprintf(stderr, "ds4: distributed worker: protocol error: %s\n", err); |
| 7783 | loop_rc = 1; |
| 7784 | break; |
| 7785 | } |
| 7786 | if (type == DS4_DIST_MSG_ERROR) { |
| 7787 | char msg[512]; |
| 7788 | uint32_t n = bytes < sizeof(msg) - 1u ? bytes : (uint32_t)sizeof(msg) - 1u; |
| 7789 | rc = dist_read_full(fd, msg, n); |
| 7790 | if (rc <= 0) { |
| 7791 | loop_rc = 1; |
no test coverage detected