Wakeup actually working (probably), add delay calculator.
This commit is contained in:
@@ -102,6 +102,7 @@ int main(int argc, char *argv[]) {
|
|||||||
hex_init();
|
hex_init();
|
||||||
|
|
||||||
peer_init();
|
peer_init();
|
||||||
|
retry_init();
|
||||||
wakeup_init();
|
wakeup_init();
|
||||||
|
|
||||||
send_init();
|
send_init();
|
||||||
@@ -116,6 +117,7 @@ int main(int argc, char *argv[]) {
|
|||||||
|
|
||||||
peer_loop();
|
peer_loop();
|
||||||
|
|
||||||
|
retry_cleanup();
|
||||||
wakeup_cleanup();
|
wakeup_cleanup();
|
||||||
send_cleanup();
|
send_cleanup();
|
||||||
|
|
||||||
|
|||||||
@@ -216,3 +216,34 @@ void uuid_gen(char *out) {
|
|||||||
uuid_generate(uuid);
|
uuid_generate(uuid);
|
||||||
uuid_unparse(uuid, out);
|
uuid_unparse(uuid, out);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
static int retry_rand_fd;
|
||||||
|
|
||||||
|
void retry_init() {
|
||||||
|
retry_rand_fd = open("/dev/urandom", O_RDONLY);
|
||||||
|
assert(retry_rand_fd >= 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
void retry_cleanup() {
|
||||||
|
assert(!close(retry_rand_fd));
|
||||||
|
}
|
||||||
|
|
||||||
|
#define RETRY_MIN_MS 2000
|
||||||
|
#define RETRY_MAX_MS 64000
|
||||||
|
#define RETRY_MULT 2
|
||||||
|
#define RETRY_MAX_JITTER_DIV 2
|
||||||
|
uint32_t retry_get_delay_ms(uint32_t prev_delay) {
|
||||||
|
uint32_t delay = prev_delay * RETRY_MULT;
|
||||||
|
delay = delay < RETRY_MIN_MS ? RETRY_MIN_MS : delay;
|
||||||
|
delay = delay > RETRY_MAX_MS ? RETRY_MAX_MS : delay;
|
||||||
|
|
||||||
|
uint32_t max_jitter = delay / RETRY_MAX_JITTER_DIV;
|
||||||
|
uint32_t jitter;
|
||||||
|
assert(read(retry_rand_fd, &jitter, sizeof(jitter)) == sizeof(jitter));
|
||||||
|
delay += jitter % max_jitter;
|
||||||
|
|
||||||
|
delay = delay > RETRY_MAX_MS ? RETRY_MAX_MS : delay;
|
||||||
|
|
||||||
|
return delay;
|
||||||
|
}
|
||||||
|
|||||||
@@ -91,3 +91,10 @@ void hex_from_int(char *, uint64_t, size_t);
|
|||||||
|
|
||||||
#define UUID_LEN 37
|
#define UUID_LEN 37
|
||||||
void uuid_gen(char *);
|
void uuid_gen(char *);
|
||||||
|
|
||||||
|
|
||||||
|
///////// retry timing
|
||||||
|
|
||||||
|
void retry_init();
|
||||||
|
void retry_cleanup();
|
||||||
|
uint32_t retry_get_delay_ms(uint32_t);
|
||||||
|
|||||||
@@ -13,7 +13,6 @@
|
|||||||
#include "common.h"
|
#include "common.h"
|
||||||
#include "incoming.h"
|
#include "incoming.h"
|
||||||
|
|
||||||
struct incoming;
|
|
||||||
struct incoming {
|
struct incoming {
|
||||||
struct peer peer;
|
struct peer peer;
|
||||||
char id[UUID_LEN];
|
char id[UUID_LEN];
|
||||||
|
|||||||
@@ -11,7 +11,6 @@
|
|||||||
#include "common.h"
|
#include "common.h"
|
||||||
#include "outgoing.h"
|
#include "outgoing.h"
|
||||||
|
|
||||||
struct outgoing;
|
|
||||||
struct outgoing {
|
struct outgoing {
|
||||||
struct peer peer;
|
struct peer peer;
|
||||||
char id[UUID_LEN];
|
char id[UUID_LEN];
|
||||||
|
|||||||
117
adsbus/wakeup.c
117
adsbus/wakeup.c
@@ -1,7 +1,13 @@
|
|||||||
|
#define _GNU_SOURCE
|
||||||
|
|
||||||
|
#include <stdlib.h>
|
||||||
#include <stdio.h>
|
#include <stdio.h>
|
||||||
#include <assert.h>
|
#include <assert.h>
|
||||||
|
#include <fcntl.h>
|
||||||
#include <unistd.h>
|
#include <unistd.h>
|
||||||
#include <errno.h>
|
#include <errno.h>
|
||||||
|
#include <time.h>
|
||||||
|
#include <string.h>
|
||||||
#include <sys/epoll.h>
|
#include <sys/epoll.h>
|
||||||
#include <pthread.h>
|
#include <pthread.h>
|
||||||
|
|
||||||
@@ -9,9 +15,53 @@
|
|||||||
|
|
||||||
#include "wakeup.h"
|
#include "wakeup.h"
|
||||||
|
|
||||||
|
struct wakeup_peer {
|
||||||
|
struct peer peer;
|
||||||
|
struct peer *inner_peer;
|
||||||
|
};
|
||||||
|
|
||||||
|
struct wakeup_request {
|
||||||
|
int fd;
|
||||||
|
uint64_t absolute_time_ms;
|
||||||
|
};
|
||||||
|
|
||||||
|
struct wakeup_entry {
|
||||||
|
struct wakeup_request request;
|
||||||
|
struct wakeup_entry *next;
|
||||||
|
};
|
||||||
|
|
||||||
static pthread_t wakeup_thread;
|
static pthread_t wakeup_thread;
|
||||||
static int wakeup_write_fd;
|
static int wakeup_write_fd;
|
||||||
|
|
||||||
|
static uint64_t wakeup_get_time_ms() {
|
||||||
|
struct timespec tp;
|
||||||
|
assert(!clock_gettime(CLOCK_MONOTONIC_COARSE, &tp));
|
||||||
|
#define NS_PER_MS 1000000
|
||||||
|
return tp.tv_sec + (tp.tv_nsec / NS_PER_MS);
|
||||||
|
}
|
||||||
|
|
||||||
|
static void wakeup_request_add(struct wakeup_entry **head, struct wakeup_request *request) {
|
||||||
|
struct wakeup_entry *entry = malloc(sizeof(*entry));
|
||||||
|
memcpy(&entry->request, request, sizeof(entry->request));
|
||||||
|
|
||||||
|
struct wakeup_entry *prev = NULL, *iter = *head;
|
||||||
|
while (iter) {
|
||||||
|
if (iter->request.absolute_time_ms > entry->request.absolute_time_ms) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
prev = iter;
|
||||||
|
iter = iter->next;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (prev) {
|
||||||
|
entry->next = prev->next;
|
||||||
|
prev->next = entry;
|
||||||
|
} else {
|
||||||
|
entry->next = *head;
|
||||||
|
*head = entry;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
static void *wakeup_main(void *arg) {
|
static void *wakeup_main(void *arg) {
|
||||||
int read_fd = (intptr_t) arg;
|
int read_fd = (intptr_t) arg;
|
||||||
|
|
||||||
@@ -20,18 +70,50 @@ static void *wakeup_main(void *arg) {
|
|||||||
|
|
||||||
struct epoll_event ev = {
|
struct epoll_event ev = {
|
||||||
.events = EPOLLIN,
|
.events = EPOLLIN,
|
||||||
|
.data = {
|
||||||
|
.fd = read_fd,
|
||||||
|
},
|
||||||
};
|
};
|
||||||
assert(!epoll_ctl(epoll_fd, EPOLL_CTL_ADD, read_fd, &ev));
|
assert(!epoll_ctl(epoll_fd, EPOLL_CTL_ADD, read_fd, &ev));
|
||||||
|
|
||||||
|
struct wakeup_entry *head = NULL;
|
||||||
|
|
||||||
while (1) {
|
while (1) {
|
||||||
#define MAX_EVENTS 10
|
#define MAX_EVENTS 1
|
||||||
struct epoll_event events[MAX_EVENTS];
|
struct epoll_event events[MAX_EVENTS];
|
||||||
int nfds = epoll_wait(epoll_fd, events, MAX_EVENTS, -1);
|
int nfds = epoll_wait(epoll_fd, events, MAX_EVENTS, -1);
|
||||||
if (nfds == -1 && errno == EINTR) {
|
if (nfds == -1 && errno == EINTR) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
assert(nfds >= 0);
|
|
||||||
break; // XXX
|
if (nfds == 1) {
|
||||||
|
assert(events[0].data.fd == read_fd);
|
||||||
|
struct wakeup_request request;
|
||||||
|
ssize_t result = read(read_fd, &request, sizeof(request));
|
||||||
|
if (result == 0) {
|
||||||
|
// Peer closed connection, shutdown thread
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
assert(result == sizeof(request));
|
||||||
|
wakeup_request_add(&head, &request);
|
||||||
|
} else {
|
||||||
|
assert(nfds == 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
uint64_t now = wakeup_get_time_ms();
|
||||||
|
while (head && head->request.absolute_time_ms < now) {
|
||||||
|
close(head->request.fd);
|
||||||
|
struct wakeup_entry *next = head->next;
|
||||||
|
free(head);
|
||||||
|
head = next;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
while (head) {
|
||||||
|
close(head->request.fd);
|
||||||
|
struct wakeup_entry *next = head->next;
|
||||||
|
free(head);
|
||||||
|
head = next;
|
||||||
}
|
}
|
||||||
|
|
||||||
assert(!close(read_fd));
|
assert(!close(read_fd));
|
||||||
@@ -39,9 +121,18 @@ static void *wakeup_main(void *arg) {
|
|||||||
return NULL;
|
return NULL;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
static void wakeup_handler(struct peer *peer) {
|
||||||
|
struct wakeup_peer *outer_peer = (struct wakeup_peer *) peer;
|
||||||
|
assert(!close(outer_peer->peer.fd));
|
||||||
|
|
||||||
|
struct peer *inner_peer = outer_peer->inner_peer;
|
||||||
|
free(outer_peer);
|
||||||
|
inner_peer->event_handler(inner_peer);
|
||||||
|
}
|
||||||
|
|
||||||
void wakeup_init() {
|
void wakeup_init() {
|
||||||
int pipefd[2];
|
int pipefd[2];
|
||||||
assert(!pipe(pipefd));
|
assert(!pipe2(pipefd, O_NONBLOCK));
|
||||||
assert(!pthread_create(&wakeup_thread, NULL, wakeup_main, (void *) (intptr_t) pipefd[0]));
|
assert(!pthread_create(&wakeup_thread, NULL, wakeup_main, (void *) (intptr_t) pipefd[0]));
|
||||||
wakeup_write_fd = pipefd[1];
|
wakeup_write_fd = pipefd[1];
|
||||||
}
|
}
|
||||||
@@ -51,5 +142,21 @@ void wakeup_cleanup() {
|
|||||||
assert(!pthread_join(wakeup_thread, NULL));
|
assert(!pthread_join(wakeup_thread, NULL));
|
||||||
}
|
}
|
||||||
|
|
||||||
void wakeup_add(struct peer *peer, int delay_ms) {
|
void wakeup_add(struct peer *peer, uint32_t delay_ms) {
|
||||||
|
int pipefd[2];
|
||||||
|
assert(!pipe2(pipefd, O_NONBLOCK));
|
||||||
|
|
||||||
|
struct wakeup_request request = {
|
||||||
|
.fd = pipefd[1],
|
||||||
|
.absolute_time_ms = wakeup_get_time_ms() + delay_ms,
|
||||||
|
};
|
||||||
|
assert(write(wakeup_write_fd, &request, sizeof(request)) == sizeof(request));
|
||||||
|
|
||||||
|
struct wakeup_peer *outer_peer = malloc(sizeof(*outer_peer));
|
||||||
|
assert(outer_peer);
|
||||||
|
outer_peer->peer.fd = pipefd[0];
|
||||||
|
outer_peer->peer.event_handler = wakeup_handler;
|
||||||
|
outer_peer->inner_peer = peer;
|
||||||
|
|
||||||
|
peer_epoll_add((struct peer *) outer_peer, EPOLLRDHUP);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,4 +4,4 @@ struct peer;
|
|||||||
|
|
||||||
void wakeup_init();
|
void wakeup_init();
|
||||||
void wakeup_cleanup();
|
void wakeup_cleanup();
|
||||||
void wakeup_add(struct peer *, int);
|
void wakeup_add(struct peer *, uint32_t);
|
||||||
|
|||||||
Reference in New Issue
Block a user