Repository navigation
Expand file tree
/
Copy pathqueue.c
More file actions
91 lines (77 loc) · 2.87 KB
/
Copy pathqueue.c
File metadata and controls
91 lines (77 loc) · 2.87 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
#include <event2/buffer.h>
#include <event2/event.h>
#include <event2/bufferevent.h>
#include "queue.h"
#include "log.h"
outgoing_queue_t *create_outgoing_queue(struct event_base *base, struct bufferevent *bev) {
outgoing_queue_t *q = malloc(sizeof(outgoing_queue_t));
q->packet_queue = evbuffer_new();
q->bev = bev;
q->min_delay_ms = 10;
q->max_delay_ms = 600;
q->timer = evtimer_new(base, process_delayed_queue, q);
return q;
}
void destroy_outgoing_queue(outgoing_queue_t *q) {
if (q == NULL) return;
if (q->timer) {
event_free(q->timer);
}
if (q->packet_queue) {
evbuffer_free(q->packet_queue);
}
free(q);
}
void outgoing_queue_add_packet(outgoing_queue_t *q, unsigned char *data, size_t size) {
evbuffer_add(q->packet_queue, &size, sizeof(size));
evbuffer_add(q->packet_queue, data, size);
if (!evtimer_pending(q->timer, NULL)) {
schedule_next_batch(q);
}
}
static void schedule_next_batch(outgoing_queue_t *q) {
int delay_ms = q->min_delay_ms + (rand() % (q->max_delay_ms - q->min_delay_ms));
struct timeval tv = {
.tv_sec = delay_ms / 1000,
.tv_usec = (delay_ms % 1000) * 1000
};
evtimer_add(q->timer, &tv);
}
static void process_delayed_queue(evutil_socket_t fd, short events, void *arg) {
outgoing_queue_t *q = (outgoing_queue_t *) arg;
size_t len = evbuffer_get_length(q->packet_queue);
if (len == 0) {
// nothing to do
return;
}
size_t packet_size;
if (evbuffer_remove(q->packet_queue, &packet_size, sizeof(packet_size)) != sizeof(packet_size)) {
log_error("process_delayed_queue(): failed to read packet_size from the outgoing packet queue");
return;
}
if (evbuffer_get_length(q->packet_queue) < packet_size) {
log_error("process_delayed_queue(): incomplete packet in queue, expected %zu bytes, got %zu", packet_size, evbuffer_get_length(q->packet_queue));
return;
}
char *packet_data = malloc(packet_size);
if (packet_data == NULL) {
log_error("process_delayed_queue(): failed to allocate memory for packet data");
return;
}
if (evbuffer_remove(q->packet_queue, packet_data, packet_size) != packet_size) {
log_error("process_delayed_queue(): failed to read packet data from queue");
free(packet_data);
return;
}
int write_result = bufferevent_write(q->bev, packet_data, packet_size);
if (write_result != 0) {
log_error("process_delayed_queue(): bufferevent_write() failed");
}
bufferevent_flush(q->bev, EV_WRITE, BEV_FLUSH);
free(packet_data);
log_debug("process_delayed_queue(): sent packet from queue, size: %zu, remaining in queue: %zu",
packet_size, evbuffer_get_length(q->packet_queue));
if (evbuffer_get_length(q->packet_queue) > 0) {
schedule_next_batch(q);
}
}