uint64_t sequence; //!< sequence number for this channel.
uint64_t ack; //!< sequence number of the other end
+ uint64_t sequence_at_last_signal; //!< when we last signaled
+
fr_time_t last_write; //!< last write to the channel
fr_time_t last_read_other; //!< last time we successfully read a message from the other the channel
fr_time_t message_interval; //!< interval between messages
}
-static int fr_channel_kevent_signal(int kq, fr_channel_control_t *cc)
+static int fr_channel_kevent_signal(fr_channel_end_t *end, fr_channel_control_t *cc)
{
struct kevent kev;
EV_SET(&kev, cc->signal, EVFILT_USER, 0, NOTE_TRIGGER | NOTE_FFCOPY, cc->ack, cc->ch);
- return kevent(kq, &kev, 1, NULL, 0, NULL);
+ end->sequence_at_last_signal = end->sequence;
+
+ return kevent(end->kq, &kev, 1, NULL, 0, NULL);
}
cc.ack = 0;
cc.ch = ch;
- return fr_channel_kevent_signal(end->kq, &cc);
+ return fr_channel_kevent_signal(end, &cc);
}
#define IALPHA (8)
cc.ack = end->ack;
cc.ch = ch;
- return fr_channel_kevent_signal(end->kq, &cc);
+ return fr_channel_kevent_signal(end, &cc);
}
*/
ack = (uint64_t) kev->data;
end = &ch->end[TO_WORKER];
- if (ack != end->sequence) {
+ if ((ack < end->sequence) &&
+ (end->sequence_at_last_signal != end->sequence)) {
end->num_resignals++;
(void) fr_channel_data_ready(ch, when, end, FR_CHANNEL_SIGNAL_DATA_TO_WORKER);
}
cc.ack = TO_WORKER;
cc.ch = ch;
- return fr_channel_kevent_signal(ch->end[TO_WORKER].kq, &cc);
+ return fr_channel_kevent_signal(&ch->end[TO_WORKER], &cc);
}
/** Acknowledge that the channel is closing
cc.ack = FROM_WORKER;
cc.ch = ch;
- return fr_channel_kevent_signal(ch->end[FROM_WORKER].kq, &cc);
+ return fr_channel_kevent_signal(&ch->end[FROM_WORKER], &cc);
}
/** Send a channel to a KQ
cc.ack = 0;
cc.ch = ch;
- return fr_channel_kevent_signal(ch->end[TO_WORKER].kq, &cc);
+ return fr_channel_kevent_signal(&ch->end[TO_WORKER], &cc);
}
void fr_channel_debug(fr_channel_t *ch, FILE *fp)