2023-02-08 23:00:04 +08:00
|
|
|
// Copyright (c) 2020-2023 Cesanta Software Limited
|
2020-12-05 19:26:32 +08:00
|
|
|
// All rights reserved
|
|
|
|
//
|
|
|
|
// Multithreading example.
|
|
|
|
// For each incoming request, we spawn a separate thread, that sleeps for
|
|
|
|
// some time to simulate long processing time, produces an output and
|
|
|
|
// hands over that output to the request handler function.
|
|
|
|
//
|
2023-02-08 23:00:04 +08:00
|
|
|
// We pass POST body to the worker thread, and respond with a calculated CRC
|
2020-12-05 19:26:32 +08:00
|
|
|
|
|
|
|
#include "mongoose.h"
|
|
|
|
|
2023-11-25 16:28:39 +08:00
|
|
|
#define LISTENING_ADDR "http://localhost:8000"
|
|
|
|
|
2023-02-08 23:00:04 +08:00
|
|
|
struct thread_data {
|
|
|
|
struct mg_queue queue; // Worker -> Connection queue
|
|
|
|
struct mg_str body; // Copy of message body
|
|
|
|
};
|
|
|
|
|
2023-11-25 16:28:39 +08:00
|
|
|
// These two helper UDP connections are used to wake up mongoose thread
|
|
|
|
static struct mg_connection *s_wakeup_server;
|
|
|
|
static struct mg_connection *s_wakeup_client;
|
|
|
|
#define WAKEUP_URL "udp://127.0.0.1:40111"
|
|
|
|
|
2023-02-13 04:30:18 +08:00
|
|
|
static void start_thread(void *(*f)(void *), void *p) {
|
2021-01-21 17:12:49 +08:00
|
|
|
#ifdef _WIN32
|
2023-02-08 23:00:04 +08:00
|
|
|
#define usleep(x) Sleep((x) / 1000)
|
2020-12-05 19:26:32 +08:00
|
|
|
_beginthread((void(__cdecl *)(void *)) f, 0, p);
|
|
|
|
#else
|
|
|
|
#include <pthread.h>
|
|
|
|
pthread_t thread_id = (pthread_t) 0;
|
|
|
|
pthread_attr_t attr;
|
|
|
|
(void) pthread_attr_init(&attr);
|
|
|
|
(void) pthread_attr_setdetachstate(&attr, PTHREAD_CREATE_DETACHED);
|
2023-02-13 04:30:18 +08:00
|
|
|
pthread_create(&thread_id, &attr, f, p);
|
2020-12-05 19:26:32 +08:00
|
|
|
pthread_attr_destroy(&attr);
|
|
|
|
#endif
|
|
|
|
}
|
|
|
|
|
2023-02-13 04:30:18 +08:00
|
|
|
static void *worker_thread(void *param) {
|
2023-02-08 23:00:04 +08:00
|
|
|
struct thread_data *d = (struct thread_data *) param;
|
|
|
|
char buf[100]; // On-stack buffer for the message queue
|
2022-08-23 23:02:30 +08:00
|
|
|
|
2023-02-13 04:30:18 +08:00
|
|
|
mg_queue_init(&d->queue, buf, sizeof(buf)); // Init queue
|
|
|
|
usleep(1 * 1000 * 1000); // Simulate long execution time
|
2022-04-22 21:42:07 +08:00
|
|
|
|
2023-02-08 23:00:04 +08:00
|
|
|
// Send a response to the connection
|
|
|
|
if (d->body.len == 0) {
|
|
|
|
mg_queue_printf(&d->queue, "Send me POST data");
|
|
|
|
} else {
|
|
|
|
uint32_t crc = mg_crc32(0, d->body.ptr, d->body.len);
|
|
|
|
mg_queue_printf(&d->queue, "crc32: %#x", crc);
|
|
|
|
free((char *) d->body.ptr);
|
2022-04-22 21:42:07 +08:00
|
|
|
}
|
2023-02-08 23:00:04 +08:00
|
|
|
|
2023-11-25 16:28:39 +08:00
|
|
|
// Wake up Mongoose thread by sending something to one of its connections
|
|
|
|
mg_send(s_wakeup_client, "hi", 2);
|
|
|
|
|
2023-02-08 23:00:04 +08:00
|
|
|
// Wait until connection reads our message, then it is safe to quit
|
2023-02-13 04:30:18 +08:00
|
|
|
while (d->queue.tail != d->queue.head) usleep(1000);
|
2023-02-08 23:00:04 +08:00
|
|
|
MG_INFO(("done, cleaning up..."));
|
|
|
|
free(d);
|
2023-02-13 04:30:18 +08:00
|
|
|
return NULL;
|
2021-01-21 17:12:49 +08:00
|
|
|
}
|
|
|
|
|
2023-11-25 16:28:39 +08:00
|
|
|
static void wfn(struct mg_connection *c, int ev, void *ev_data, void *fn_data) {
|
|
|
|
if (ev == MG_EV_READ) {
|
|
|
|
c->recv.len = 0; // Discard received data
|
|
|
|
}
|
|
|
|
(void) ev_data, (void) fn_data;
|
|
|
|
}
|
|
|
|
|
2020-12-05 19:26:32 +08:00
|
|
|
// HTTP request callback
|
2021-08-12 02:17:04 +08:00
|
|
|
static void fn(struct mg_connection *c, int ev, void *ev_data, void *fn_data) {
|
2020-12-05 19:26:32 +08:00
|
|
|
if (ev == MG_EV_HTTP_MSG) {
|
2023-02-08 23:00:04 +08:00
|
|
|
// Received HTTP request. Allocate thread data and spawn a worker thread
|
2020-12-05 19:26:32 +08:00
|
|
|
struct mg_http_message *hm = (struct mg_http_message *) ev_data;
|
2023-02-08 23:00:04 +08:00
|
|
|
struct thread_data *d = (struct thread_data *) calloc(1, sizeof(*d));
|
|
|
|
d->body = mg_strdup(hm->body); // Pass received body to the worker
|
|
|
|
start_thread(worker_thread, d); // Start a thread
|
|
|
|
*(void **) c->data = d; // Memorise data pointer in c->data
|
|
|
|
} else if (ev == MG_EV_POLL) {
|
|
|
|
// Poll event. Delivered to us every mg_mgr_poll interval or faster
|
|
|
|
struct thread_data *d = *(struct thread_data **) c->data;
|
|
|
|
size_t len;
|
|
|
|
char *buf;
|
|
|
|
// Check if we have a message from the worker
|
2023-02-13 04:30:18 +08:00
|
|
|
if (d != NULL && (len = mg_queue_next(&d->queue, &buf)) > 0) {
|
2023-02-08 23:00:04 +08:00
|
|
|
// Got message from worker. Send a response and cleanup
|
|
|
|
mg_http_reply(c, 200, "", "%.*s\n", (int) len, buf);
|
2023-02-13 04:30:18 +08:00
|
|
|
mg_queue_del(&d->queue, len); // Delete message
|
|
|
|
*(void **) c->data = NULL; // Forget about thread data
|
2020-12-05 19:26:32 +08:00
|
|
|
}
|
|
|
|
}
|
2023-02-08 23:00:04 +08:00
|
|
|
(void) fn_data;
|
2020-12-05 19:26:32 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
int main(void) {
|
|
|
|
struct mg_mgr mgr;
|
|
|
|
mg_mgr_init(&mgr);
|
2022-08-01 18:19:32 +08:00
|
|
|
mg_log_set(MG_LL_DEBUG); // Set debug log level
|
2023-11-25 16:28:39 +08:00
|
|
|
mg_http_listen(&mgr, LISTENING_ADDR, fn, NULL);
|
|
|
|
s_wakeup_server = mg_listen(&mgr, WAKEUP_URL, wfn, NULL);
|
|
|
|
s_wakeup_client = mg_connect(&mgr, WAKEUP_URL, wfn, NULL);
|
|
|
|
for (;;) mg_mgr_poll(&mgr, 5000); // Event loop. Use 5s poll interval
|
|
|
|
mg_mgr_free(&mgr); // Cleanup
|
2020-12-05 19:26:32 +08:00
|
|
|
return 0;
|
|
|
|
}
|