*/
check_close:
if (!signaled_close && (num_messages >= max_messages) && (num_outstanding == 0)) {
- rcode = fr_channel_signal_close(channel, false);
+ rcode = fr_channel_signal_worker_close(channel);
if (rcode < 0) {
fprintf(stderr, "Failed signaling close: %s\n", strerror(errno));
exit(1);
fr_message_done(&cd->m);
}
- (void) fr_channel_signal_close(channel, true);
+ (void) fr_channel_ack_worker_close(channel);
continue;
}
return ch->active;
}
-/** Signal a channel that it is closing.
+/** Signal a worker that the channel is closing
*
* @param[in] ch The channel.
- * @param[in] ack Whether we're acking a previous request to close the channel.
* @return
* - <0 on error
* - 0 on success
*/
-int fr_channel_signal_close(fr_channel_t *ch, bool ack)
+int fr_channel_signal_worker_close(fr_channel_t *ch)
{
fr_channel_control_t cc;
ch->active = false;
cc.signal = FR_CHANNEL_SIGNAL_CLOSE;
- cc.ack = ack;
+ cc.ack = TO_WORKER;
cc.ch = ch;
- return fr_channel_kevent_signal(ch->end[ack].kq, &cc);
+ return fr_channel_kevent_signal(ch->end[TO_WORKER].kq, &cc);
+}
+
+/** Acknowledge that the channel is closing
+ *
+ * @param[in] ch The channel.
+ * @return
+ * - <0 on error
+ * - 0 on success
+ */
+int fr_channel_ack_worker_close(fr_channel_t *ch)
+{
+ fr_channel_control_t cc;
+
+ ch->active = false;
+
+ cc.signal = FR_CHANNEL_SIGNAL_CLOSE;
+ cc.ack = FROM_WORKER;
+ cc.ch = ch;
+
+ return fr_channel_kevent_signal(ch->end[FROM_WORKER].kq, &cc);
}
/** Send a channel to a KQ
bool fr_channel_active(fr_channel_t *ch) CC_HINT(nonnull);
int fr_channel_signal_open(int kq, fr_channel_t *ch) CC_HINT(nonnull);
-int fr_channel_signal_close(fr_channel_t *ch, bool ack) CC_HINT(nonnull);
+int fr_channel_signal_worker_close(fr_channel_t *ch) CC_HINT(nonnull);
+int fr_channel_ack_worker_close(fr_channel_t *ch) CC_HINT(nonnull);
void fr_channel_debug(fr_channel_t *ch, FILE *fp);
* closing/
*/
while ((worker = fr_heap_pop(rc->workers)) != NULL) {
- fr_channel_signal_close(worker->channel, false);
+ fr_channel_signal_worker_close(worker->channel);
(void) fr_heap_insert(rc->closing, worker);
}
* automatically freed when our talloc context is freed.
*/
for (i = 0; i < worker->num_channels; i++) {
- fr_channel_signal_close(worker->channel[i], true);
+ fr_channel_ack_worker_close(worker->channel[i]);
}
}