move to nginx-like worker eventloop
This commit is contained in:
@@ -5,14 +5,12 @@
|
||||
#include <fcntl.h>
|
||||
#include <netinet/in.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/epoll.h>
|
||||
#include <sys/types.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include "http/http.h"
|
||||
#include "server/threadpool.h"
|
||||
#include "server/worker.h"
|
||||
#include "utils/colors.h"
|
||||
#include "utils/utils.h"
|
||||
|
||||
@@ -83,130 +81,41 @@ cws_server_ret cws_server_start(cws_config *config) {
|
||||
return CWS_SERVER_OK;
|
||||
}
|
||||
|
||||
static cws_server_ret cws_server_setup_epoll(int sockfd, int *epfd_out) {
|
||||
int epfd = epoll_create1(0);
|
||||
if (epfd == -1) {
|
||||
CWS_LOG_ERROR("epoll_create(): %s", strerror(errno));
|
||||
|
||||
return CWS_SERVER_EPOLL_CREATE_ERROR;
|
||||
}
|
||||
|
||||
cws_server_ret ret;
|
||||
|
||||
ret = cws_fd_set_nonblocking(sockfd);
|
||||
if (ret != CWS_SERVER_OK) {
|
||||
close(epfd);
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
ret = cws_epoll_add(epfd, sockfd, EPOLLIN | EPOLLET);
|
||||
if (ret != CWS_SERVER_OK) {
|
||||
close(epfd);
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
*epfd_out = epfd;
|
||||
|
||||
return CWS_SERVER_OK;
|
||||
}
|
||||
|
||||
cws_server_ret cws_server_loop(int server_fd, cws_config *config) {
|
||||
cws_server_ret ret;
|
||||
int epfd;
|
||||
|
||||
ret = cws_server_setup_epoll(server_fd, &epfd);
|
||||
if (ret != CWS_SERVER_OK) {
|
||||
return ret;
|
||||
}
|
||||
|
||||
mcl_hashmap *clients = mcl_hm_init(my_int_hash_fn, my_int_equal_fn, my_int_free_key_fn, my_str_free_fn, sizeof(int), sizeof(char) * INET_ADDRSTRLEN);
|
||||
if (!clients) {
|
||||
if (clients == NULL) {
|
||||
return CWS_SERVER_HASHMAP_INIT;
|
||||
}
|
||||
|
||||
struct epoll_event *revents = malloc(CWS_SERVER_EPOLL_MAXEVENTS * sizeof(struct epoll_event));
|
||||
if (!revents) {
|
||||
size_t workers_num = 4;
|
||||
size_t workers_index = 0;
|
||||
cws_worker **workers = cws_worker_init(workers_num, clients, config);
|
||||
if (workers == NULL) {
|
||||
mcl_hm_free(clients);
|
||||
close(epfd);
|
||||
|
||||
return CWS_SERVER_MALLOC_ERROR;
|
||||
}
|
||||
|
||||
cws_threadpool *pool = cws_threadpool_init(4, 16);
|
||||
if (pool == NULL) {
|
||||
mcl_hm_free(clients);
|
||||
close(epfd);
|
||||
|
||||
return CWS_SERVER_THREADPOOL_ERROR;
|
||||
return CWS_SERVER_WORKER_ERROR;
|
||||
}
|
||||
|
||||
int client_fd = 0;
|
||||
while (cws_server_run) {
|
||||
int nfds = epoll_wait(epfd, revents, CWS_SERVER_EPOLL_MAXEVENTS, CWS_SERVER_EPOLL_TIMEOUT);
|
||||
|
||||
if (nfds == 0) {
|
||||
/* TODO: Check for inactive clients */
|
||||
client_fd = cws_server_handle_new_client(server_fd, clients);
|
||||
if (client_fd < 0) {
|
||||
continue;
|
||||
}
|
||||
|
||||
for (int i = 0; i < nfds; ++i) {
|
||||
if (revents[i].data.fd == server_fd) {
|
||||
/* Handle new client */
|
||||
int client_fd = cws_server_handle_new_client(server_fd, epfd, clients);
|
||||
if (client_fd < 0) {
|
||||
continue;
|
||||
}
|
||||
cws_epoll_add(epfd, client_fd, EPOLLIN | EPOLLET);
|
||||
cws_fd_set_nonblocking(client_fd);
|
||||
mcl_hm_set(clients, &client_fd, "test");
|
||||
} else {
|
||||
/* Handle client data */
|
||||
int client_fd = revents[i].data.fd;
|
||||
|
||||
cws_task *task = malloc(sizeof(cws_task));
|
||||
if (!task) {
|
||||
cws_server_close_client(epfd, client_fd, clients);
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
task->client_fd = client_fd;
|
||||
task->config = config;
|
||||
task->epfd = epfd;
|
||||
task->clients = clients;
|
||||
|
||||
cws_thread_task *ttask = malloc(sizeof(cws_thread_task));
|
||||
if (!ttask) {
|
||||
free(task);
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
ttask->arg = task;
|
||||
ttask->function = (void *)cws_server_handle_client_data;
|
||||
|
||||
ret = cws_threadpool_submit(pool, ttask);
|
||||
if (ret != CWS_SERVER_OK) {
|
||||
cws_server_close_client(epfd, client_fd, clients);
|
||||
free(task);
|
||||
}
|
||||
|
||||
free(ttask);
|
||||
}
|
||||
}
|
||||
/* Add client to worker */
|
||||
int random = 10;
|
||||
mcl_hm_set(clients, &client_fd, &random);
|
||||
write(workers[workers_index]->pipefd[1], &client_fd, sizeof(int));
|
||||
workers_index = (workers_index + 1) % workers_num;
|
||||
}
|
||||
|
||||
/* Clean up everything */
|
||||
free(revents);
|
||||
close(epfd);
|
||||
cws_worker_free(workers, workers_num);
|
||||
mcl_hm_free(clients);
|
||||
cws_threadpool_destroy(pool);
|
||||
|
||||
return CWS_SERVER_OK;
|
||||
}
|
||||
|
||||
int cws_server_handle_new_client(int server_fd, int epfd, mcl_hashmap *clients) {
|
||||
int cws_server_handle_new_client(int server_fd, mcl_hashmap *clients) {
|
||||
struct sockaddr_storage their_sa;
|
||||
socklen_t theirsa_size = sizeof their_sa;
|
||||
char ip[INET_ADDRSTRLEN];
|
||||
@@ -222,106 +131,16 @@ int cws_server_handle_new_client(int server_fd, int epfd, mcl_hashmap *clients)
|
||||
return client_fd;
|
||||
}
|
||||
|
||||
static size_t cws_read_data(int sockfd, mcl_string **str) {
|
||||
size_t total_bytes = 0;
|
||||
ssize_t bytes_read;
|
||||
int cws_server_accept_client(int server_fd, struct sockaddr_storage *their_sa, socklen_t *theirsa_size) {
|
||||
const int client_fd = accept(server_fd, (struct sockaddr *)their_sa, theirsa_size);
|
||||
|
||||
if (*str == NULL) {
|
||||
*str = mcl_string_new("", 4096);
|
||||
}
|
||||
|
||||
char tmp[4096];
|
||||
memset(tmp, 0, sizeof(tmp));
|
||||
|
||||
while (1) {
|
||||
bytes_read = recv(sockfd, tmp, sizeof(tmp), 0);
|
||||
|
||||
if (bytes_read == -1) {
|
||||
if (errno == EAGAIN || errno == EWOULDBLOCK) {
|
||||
break;
|
||||
}
|
||||
|
||||
CWS_LOG_ERROR("recv(): %s", strerror(errno));
|
||||
|
||||
return -1;
|
||||
} else if (bytes_read == 0) {
|
||||
return -1;
|
||||
if (client_fd == -1) {
|
||||
if (errno != EWOULDBLOCK) {
|
||||
CWS_LOG_ERROR("accept(): %s", strerror(errno));
|
||||
}
|
||||
|
||||
total_bytes += bytes_read;
|
||||
mcl_string_append(*str, tmp);
|
||||
}
|
||||
|
||||
return total_bytes;
|
||||
}
|
||||
|
||||
void *cws_server_handle_client_data(void *arg) {
|
||||
cws_task *task = (cws_task *)arg;
|
||||
|
||||
int client_fd = task->client_fd;
|
||||
int epfd = task->epfd;
|
||||
mcl_hashmap *clients = task->clients;
|
||||
cws_config *config = task->config;
|
||||
|
||||
/* Read data from socket */
|
||||
mcl_string *data = NULL;
|
||||
size_t total_bytes = cws_read_data(client_fd, &data);
|
||||
if (total_bytes <= 0) {
|
||||
if (data) {
|
||||
mcl_string_free(data);
|
||||
}
|
||||
cws_server_close_client(epfd, client_fd, clients);
|
||||
free(task);
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
/* Parse HTTP request */
|
||||
cws_http *request = cws_http_parse(data, client_fd, config);
|
||||
mcl_string_free(data);
|
||||
|
||||
if (request == NULL) {
|
||||
cws_server_close_client(epfd, client_fd, clients);
|
||||
free(task);
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
int keepalive = cws_http_send_resource(request);
|
||||
cws_http_free(request);
|
||||
if (!keepalive) {
|
||||
CWS_LOG_DEBUG("keep-alive: %d, closing fd: %d", keepalive, client_fd);
|
||||
cws_server_close_client(epfd, client_fd, clients);
|
||||
}
|
||||
|
||||
free(task);
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
cws_server_ret cws_epoll_add(int epfd, int sockfd, uint32_t events) {
|
||||
struct epoll_event event;
|
||||
event.events = events;
|
||||
event.data.fd = sockfd;
|
||||
const int status = epoll_ctl(epfd, EPOLL_CTL_ADD, sockfd, &event);
|
||||
|
||||
if (status != 0) {
|
||||
CWS_LOG_ERROR("epoll_ctl_add(): %s", strerror(errno));
|
||||
return CWS_SERVER_EPOLL_ADD_ERROR;
|
||||
}
|
||||
|
||||
return CWS_SERVER_OK;
|
||||
}
|
||||
|
||||
cws_server_ret cws_epoll_del(int epfd, int sockfd) {
|
||||
const int status = epoll_ctl(epfd, EPOLL_CTL_DEL, sockfd, NULL);
|
||||
|
||||
if (status != 0) {
|
||||
CWS_LOG_ERROR("epoll_ctl_del(): %s", strerror(errno));
|
||||
return CWS_SERVER_EPOLL_DEL_ERROR;
|
||||
}
|
||||
|
||||
return CWS_SERVER_OK;
|
||||
return client_fd;
|
||||
}
|
||||
|
||||
cws_server_ret cws_fd_set_nonblocking(int sockfd) {
|
||||
@@ -334,26 +153,3 @@ cws_server_ret cws_fd_set_nonblocking(int sockfd) {
|
||||
|
||||
return CWS_SERVER_OK;
|
||||
}
|
||||
|
||||
int cws_server_accept_client(int server_fd, struct sockaddr_storage *their_sa, socklen_t *theirsa_size) {
|
||||
const int client_fd = accept(server_fd, (struct sockaddr *)their_sa, theirsa_size);
|
||||
|
||||
if (client_fd == -1) {
|
||||
if (errno != EWOULDBLOCK) {
|
||||
CWS_LOG_ERROR("accept(): %s", strerror(errno));
|
||||
}
|
||||
return -1;
|
||||
}
|
||||
|
||||
return client_fd;
|
||||
}
|
||||
|
||||
void cws_server_close_client(int epfd, int client_fd, mcl_hashmap *clients) {
|
||||
/* TODO: fix race conditions */
|
||||
mcl_bucket *client = mcl_hm_get(clients, &client_fd);
|
||||
if (client) {
|
||||
cws_epoll_del(epfd, client_fd);
|
||||
mcl_hm_remove(clients, &client_fd);
|
||||
mcl_hm_free_bucket(client);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user