squeezed in raw proxying code into normal proxy

This commit is contained in:
hyung-hwan 2014-07-15 16:22:24 +00:00
parent a7ca23fa50
commit a0e2a7067c
5 changed files with 397 additions and 252 deletions

View File

@ -480,9 +480,9 @@ struct qse_httpd_rsrc_cgi_t
typedef struct qse_httpd_rsrc_proxy_t qse_httpd_rsrc_proxy_t; typedef struct qse_httpd_rsrc_proxy_t qse_httpd_rsrc_proxy_t;
struct qse_httpd_rsrc_proxy_t struct qse_httpd_rsrc_proxy_t
{ {
qse_nwad_t dst; qse_nwad_t dst; /* remote destination address to connect to */
qse_nwad_t src; qse_nwad_t src; /* local binding address */
int raw; int raw; /* raw or normal */
}; };
typedef struct qse_httpd_rsrc_dir_t qse_httpd_rsrc_dir_t; typedef struct qse_httpd_rsrc_dir_t qse_httpd_rsrc_dir_t;

View File

@ -37,8 +37,6 @@ struct task_proxy_t
int init_failed; int init_failed;
qse_httpd_t* httpd; qse_httpd_t* httpd;
const qse_mchar_t* host;
int method; int method;
qse_http_version_t version; qse_http_version_t version;
int keepalive; /* taken from the request */ int keepalive; /* taken from the request */
@ -197,6 +195,12 @@ static int proxy_snatch_client_input_raw (
task = (qse_httpd_task_t*)ctx; task = (qse_httpd_task_t*)ctx;
proxy = (task_proxy_t*)task->ctx; proxy = (task_proxy_t*)task->ctx;
/* this function is never called with ptr of QSE_NULL
* because this callback is set manually after the request
* has been discarded or completed in task_init_proxy() and
* qse_htre_completecontent or qse-htre_discardcontent() is
* not called again. Unlinkw proxy_snatch_client_input(),
* it doesn't care about EOF indicated by ptr of QSE_NULL. */
if (ptr && !(proxy->reqflags & PROXY_REQ_FWDERR)) if (ptr && !(proxy->reqflags & PROXY_REQ_FWDERR))
{ {
if (qse_mbs_ncat (proxy->reqfwdbuf, ptr, len) == (qse_size_t)-1) if (qse_mbs_ncat (proxy->reqfwdbuf, ptr, len) == (qse_size_t)-1)
@ -732,6 +736,8 @@ to the head all the time.. grow the buffer to a certain limit. */
} }
} }
/* ------------------------------------------------------------------------ */
static int task_init_proxy ( static int task_init_proxy (
qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task) qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task)
{ {
@ -776,7 +782,6 @@ static int task_init_proxy (
/* the caller must make sure that the actual content is discarded or completed /* the caller must make sure that the actual content is discarded or completed
* and the following data is treated as contents */ * and the following data is treated as contents */
printf ("proxy req = %p %d %d\n", arg->req, (arg->req->state & QSE_HTRE_DISCARDED), (arg->req->state& QSE_HTRE_COMPLETED));
QSE_ASSERT (arg->req->state & (QSE_HTRE_DISCARDED | QSE_HTRE_COMPLETED)); QSE_ASSERT (arg->req->state & (QSE_HTRE_DISCARDED | QSE_HTRE_COMPLETED));
QSE_ASSERT (qse_htrd_getoption(client->htrd) & QSE_HTRD_DUMMY); QSE_ASSERT (qse_htrd_getoption(client->htrd) & QSE_HTRD_DUMMY);
@ -946,6 +951,8 @@ qse_printf (QSE_T("GOING TO PROXY [%hs]\n"), QSE_MBS_PTR(proxy->reqfwdbuf));
return 0; return 0;
oops: oops:
printf ("init_proxy failed...........................................\n");
/* since a new task can't be added in the initializer, /* since a new task can't be added in the initializer,
* i mark that initialization failed and let task_main_proxy() * i mark that initialization failed and let task_main_proxy()
* add an error task */ * add an error task */
@ -961,11 +968,14 @@ oops:
return 0; return 0;
} }
/* ------------------------------------------------------------------------ */
static void task_fini_proxy ( static void task_fini_proxy (
qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task) qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task)
{ {
task_proxy_t* proxy = (task_proxy_t*)task->ctx; task_proxy_t* proxy = (task_proxy_t*)task->ctx;
printf ("task_fini_proxy.................\n");
if (proxy->peer_status & PROXY_PEER_OPEN) if (proxy->peer_status & PROXY_PEER_OPEN)
httpd->opt.scb.peer.close (httpd, &proxy->peer); httpd->opt.scb.peer.close (httpd, &proxy->peer);
@ -975,6 +985,8 @@ static void task_fini_proxy (
if (proxy->req) qse_htre_unsetconcb (proxy->req); if (proxy->req) qse_htre_unsetconcb (proxy->req);
} }
/* ------------------------------------------------------------------------ */
static int task_main_proxy_5 ( static int task_main_proxy_5 (
qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task) qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task)
{ {
@ -982,7 +994,7 @@ static int task_main_proxy_5 (
qse_ssize_t n; qse_ssize_t n;
#if 0 #if 0
qse_printf (QSE_T("task_main_proxy_5 trigger[0].mask=%d trigger[1].mask=%d trigger[2].mask=%d\n"), printf ("task_main_proxy_5 trigger[0].mask=%d trigger[1].mask=%d trigger[2].mask=%d\n",
task->trigger[0].mask, task->trigger[1].mask, task->trigger[2].mask); task->trigger[0].mask, task->trigger[1].mask, task->trigger[2].mask);
#endif #endif
@ -993,7 +1005,7 @@ qse_printf (QSE_T("task_main_proxy_5 trigger[0].mask=%d trigger[1].mask=%d trigg
} }
else if (task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_WRITABLE) else if (task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_WRITABLE)
{ {
/* if the peer side is writable */ /* if the peer side is writable while the client side is not readable*/
proxy_forward_client_input_to_peer (httpd, task, 1); proxy_forward_client_input_to_peer (httpd, task, 1);
} }
@ -1030,8 +1042,8 @@ static int task_main_proxy_4 (
{ {
task_proxy_t* proxy = (task_proxy_t*)task->ctx; task_proxy_t* proxy = (task_proxy_t*)task->ctx;
#if 0 #if 1
qse_printf (QSE_T("task_main_proxy_4 trigger[0].mask=%d trigger[1].mask=%d trigger[2].mask=%d\n"), printf ("task_main_proxy_4 trigger[0].mask=%d trigger[1].mask=%d trigger[2].mask=%d\n",
task->trigger[0].mask, task->trigger[1].mask, task->trigger[2].mask); task->trigger[0].mask, task->trigger[1].mask, task->trigger[2].mask);
#endif #endif
@ -1067,8 +1079,11 @@ qse_printf (QSE_T("task_main_proxy_4 trigger[0].mask=%d trigger[1].mask=%d trigg
} }
if (n == 0) if (n == 0)
{ {
/* peer closed connection */
if (proxy->resflags & PROXY_RES_PEER_LENGTH) if (proxy->resflags & PROXY_RES_PEER_LENGTH)
{ {
QSE_ASSERT (!proxy->raw);
if (proxy->peer_output_received < proxy->peer_output_length) if (proxy->peer_output_received < proxy->peer_output_length)
{ {
if (httpd->opt.trait & QSE_HTTPD_LOGACT) if (httpd->opt.trait & QSE_HTTPD_LOGACT)
@ -1078,8 +1093,23 @@ qse_printf (QSE_T("task_main_proxy_4 trigger[0].mask=%d trigger[1].mask=%d trigg
} }
task->main = task_main_proxy_5; task->main = task_main_proxy_5;
/* nothing to read from peer. set the mask to 0 */
task->trigger[0].mask = 0; task->trigger[0].mask = 0;
/* arrange to be called if the client side is writable */
task->trigger[2].mask |= QSE_HTTPD_TASK_TRIGGER_WRITE; task->trigger[2].mask |= QSE_HTTPD_TASK_TRIGGER_WRITE;
if (proxy->raw)
{
/* peer connection has been closed.
* so no more forwarding from the client to the peer
* is possible. get rid of the content callback on the
* client side. */
qse_htre_unsetconcb (proxy->req);
proxy->req = QSE_NULL;
}
return 1; return 1;
} }
@ -1088,6 +1118,8 @@ qse_printf (QSE_T("task_main_proxy_4 trigger[0].mask=%d trigger[1].mask=%d trigg
if (proxy->resflags & PROXY_RES_PEER_LENGTH) if (proxy->resflags & PROXY_RES_PEER_LENGTH)
{ {
QSE_ASSERT (!proxy->raw);
if (proxy->peer_output_received > proxy->peer_output_length) if (proxy->peer_output_received > proxy->peer_output_length)
{ {
/* proxy returning too much data... something is wrong in PROXY */ /* proxy returning too much data... something is wrong in PROXY */
@ -1208,16 +1240,18 @@ static int task_main_proxy_2 (
int http_errnum = 500; int http_errnum = 500;
#if 0 #if 0
qse_printf (QSE_T("task_main_proxy_2 trigger[0].mask=%d trigger[1].mask=%d trigger[2].mask=%d\n"), printf ("task_main_proxy_2 trigger[0].mask=%d trigger[1].mask=%d trigger[2].mask=%d\n",
task->trigger[0].mask, task->trigger[1].mask, task->trigger[2].mask); task->trigger[0].mask, task->trigger[1].mask, task->trigger[2].mask);
#endif #endif
if (task->trigger[2].mask & QSE_HTTPD_TASK_TRIGGER_READABLE) if (task->trigger[2].mask & QSE_HTTPD_TASK_TRIGGER_READABLE)
{ {
/* client is readable */
proxy_forward_client_input_to_peer (httpd, task, 0); proxy_forward_client_input_to_peer (httpd, task, 0);
} }
else if (task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_WRITABLE) else if (task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_WRITABLE)
{ {
/* client is not readable but peer is writable */
proxy_forward_client_input_to_peer (httpd, task, 1); proxy_forward_client_input_to_peer (httpd, task, 1);
} }
@ -1342,7 +1376,6 @@ for (i = 0; i < proxy->buflen; i++) qse_printf (QSE_T("%hc"), proxy->buf[i]);
qse_printf (QSE_T("]\n")); qse_printf (QSE_T("]\n"));
#endif #endif
if (qse_htrd_feed (proxy->peer_htrd, proxy->buf, proxy->buflen) <= -1) if (qse_htrd_feed (proxy->peer_htrd, proxy->buf, proxy->buflen) <= -1)
{ {
if (httpd->opt.trait & QSE_HTTPD_LOGACT) if (httpd->opt.trait & QSE_HTTPD_LOGACT)
@ -1413,7 +1446,6 @@ static int task_main_proxy_1 (
int http_errnum = 500; int http_errnum = 500;
/* wait for peer to get connected */ /* wait for peer to get connected */
if (task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_READABLE || if (task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_READABLE ||
task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_WRITABLE) task->trigger[0].mask & QSE_HTTPD_TASK_TRIGGER_WRITABLE)
{ {
@ -1438,7 +1470,6 @@ static int task_main_proxy_1 (
if (n >= 1) if (n >= 1)
{ {
/* connected to the peer now */ /* connected to the peer now */
proxy->peer_status |= PROXY_PEER_CONNECTED; proxy->peer_status |= PROXY_PEER_CONNECTED;
if (proxy->req) if (proxy->req)
@ -1461,11 +1492,27 @@ static int task_main_proxy_1 (
} }
if (proxy->raw) if (proxy->raw)
task->main = task_main_proxy_4; {
printf ("SWITCHING TO PROXY 3...%p\n", proxy->req);
if (qse_mbs_fmt (proxy->res, QSE_MT("HTTP/%d.%d 200 Connection established\r\n\r\n"),
(int)proxy->version.major, (int)proxy->version.minor) == (qse_size_t)-1)
{
proxy->httpd->errnum = QSE_HTTPD_ENOMEM;
goto oops;
}
proxy->res_pending = QSE_MBS_LEN(proxy->res) - proxy->res_consumed;
/* arrange to be called if the client side is writable.
* it must write the injected response. */
task->trigger[2].mask |= QSE_HTTPD_TASK_TRIGGER_WRITE;
task->main = task_main_proxy_3;
}
else else
{
task->main = task_main_proxy_2; task->main = task_main_proxy_2;
} }
} }
}
return 1; return 1;
@ -1544,8 +1591,29 @@ static int task_main_proxy (
task->trigger[0].mask |= QSE_HTTPD_TASK_TRIGGER_WRITE; task->trigger[0].mask |= QSE_HTTPD_TASK_TRIGGER_WRITE;
} }
} }
if (proxy->raw)
{
/* TODO: write response */
/* inject http response */
if (qse_mbs_fmt (proxy->res, QSE_MT("HTTP/%d.%d 200 Connection established\r\n\r\n"),
(int)proxy->version.major, (int)proxy->version.minor) == (qse_size_t)-1)
{
proxy->httpd->errnum = QSE_HTTPD_ENOMEM;
goto oops;
}
proxy->res_pending = QSE_MBS_LEN(proxy->res) - proxy->res_consumed;
/* arrange to be called if the client side is writable.
* it must write the injected response. */
task->trigger[2].mask |= QSE_HTTPD_TASK_TRIGGER_WRITE;
task->main = task_main_proxy_3;
}
else
{
task->main = task_main_proxy_2; task->main = task_main_proxy_2;
} }
}
return 1; return 1;
@ -1566,6 +1634,8 @@ oops:
proxy->method, &proxy->version, proxy->keepalive) == QSE_NULL)? -1: 0; proxy->method, &proxy->version, proxy->keepalive) == QSE_NULL)? -1: 0;
} }
/* ------------------------------------------------------------------------ */
qse_httpd_task_t* qse_httpd_entaskproxy ( qse_httpd_task_t* qse_httpd_entaskproxy (
qse_httpd_t* httpd, qse_httpd_t* httpd,
qse_httpd_client_t* client, qse_httpd_client_t* client,

View File

@ -29,9 +29,8 @@
typedef struct task_resol_arg_t task_resol_arg_t; typedef struct task_resol_arg_t task_resol_arg_t;
struct task_resol_arg_t struct task_resol_arg_t
{ {
const qse_mchar_t* path; const qse_mchar_t* host;
qse_htre_t* req; const qse_htre_t* req;
int nph;
}; };
typedef struct task_resol_t task_resol_t; typedef struct task_resol_t task_resol_t;
@ -40,21 +39,32 @@ struct task_resol_t
int init_failed; int init_failed;
qse_httpd_t* httpd; qse_httpd_t* httpd;
const qse_mchar_t* path; int method;
qse_http_version_t version; qse_http_version_t version;
int keepalive; /* taken from the request */ int keepalive; /* taken from the request */
qse_mchar_t* host;
}; };
static int task_init_resol ( static int task_init_resol (
qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task) qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task)
{ {
task_resol_t* resol; task_resol_t* resol;
task_resol_arg_t* arg;
resol = (task_resol_t*)qse_httpd_gettaskxtn (httpd, task); resol = (task_resol_t*)qse_httpd_gettaskxtn (httpd, task);
arg = (task_resol_arg_t*)task->ctx;
QSE_MEMSET (resol, 0, QSE_SIZEOF(*resol)); QSE_MEMSET (resol, 0, QSE_SIZEOF(*resol));
resol->httpd = httpd; resol->httpd = httpd;
resol->method = qse_htre_getqmethodtype(arg->req);
resol->version = *qse_htre_getversion(arg->req);
resol->keepalive = (arg->req->attr.flags & QSE_HTRE_ATTR_KEEPALIVE);
resol->host = (qse_mchar_t*)(resol + 1);
qse_mbscpy (resol->host, arg->host);
task->ctx = resol; task->ctx = resol;
return 0; return 0;
} }
@ -68,6 +78,11 @@ static void task_fini_resol (
static int task_main_resol ( static int task_main_resol (
qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task) qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task)
{ {
/* dns.open ();
dns.send (...);
dns.close ();*/
return 0; return 0;
} }
@ -75,11 +90,15 @@ qse_httpd_task_t* qse_httpd_entaskresol (
qse_httpd_t* httpd, qse_httpd_t* httpd,
qse_httpd_client_t* client, qse_httpd_client_t* client,
qse_httpd_task_t* pred, qse_httpd_task_t* pred,
const qse_mchar_t* host) const qse_mchar_t* host,
qse_htre_t* req)
{ {
qse_httpd_task_t task; qse_httpd_task_t task;
task_resol_arg_t arg; task_resol_arg_t arg;
arg.host = host;
arg.req = req;
QSE_MEMSET (&task, 0, QSE_SIZEOF(task)); QSE_MEMSET (&task, 0, QSE_SIZEOF(task));
task.init = task_init_resol; task.init = task_init_resol;
task.fini = task_fini_resol; task.fini = task_fini_resol;
@ -87,7 +106,7 @@ qse_httpd_task_t* qse_httpd_entaskresol (
task.ctx = &arg; task.ctx = &arg;
return qse_httpd_entask ( return qse_httpd_entask (
httpd, client, pred, &task, QSE_SIZEOF(task_resol_t) httpd, client, pred, &task, QSE_SIZEOF(task_resol_t) + qse_mbslen(host) + 1
); );
} }

View File

@ -1126,10 +1126,15 @@ static void dispatch_muxcb (qse_mux_t* mux, const qse_mux_evt_t* evt)
{ {
mux_xtn_t* xtn; mux_xtn_t* xtn;
qse_ubi_t ubi; qse_ubi_t ubi;
int mask = 0;
xtn = qse_mux_getxtn (mux); xtn = qse_mux_getxtn (mux);
ubi.i = evt->hnd; ubi.i = evt->hnd;
xtn->cbfun (xtn->httpd, mux, ubi, evt->mask, evt->data);
if (evt->mask & QSE_MUX_IN) mask |= QSE_HTTPD_MUX_READ;
if (evt->mask & QSE_MUX_OUT) mask |= QSE_HTTPD_MUX_WRITE;
xtn->cbfun (xtn->httpd, mux, ubi, mask, evt->data);
} }
static void* mux_open (qse_httpd_t* httpd, qse_httpd_muxcb_t cbfun) static void* mux_open (qse_httpd_t* httpd, qse_httpd_muxcb_t cbfun)
@ -2041,7 +2046,9 @@ if (qse_htre_getcontentlen(req) > 0)
* 'Expect: 100-continue' and 'Connection: keep-alive'. */ * 'Expect: 100-continue' and 'Connection: keep-alive'. */
qse_httpd_discardcontent (httpd, req); qse_httpd_discardcontent (httpd, req);
} }
else if (mth == QSE_HTTP_POST && else
{
if (mth == QSE_HTTP_POST &&
!(req->attr.flags & QSE_HTRE_ATTR_LENGTH) && !(req->attr.flags & QSE_HTRE_ATTR_LENGTH) &&
!(req->attr.flags & QSE_HTRE_ATTR_CHUNKED)) !(req->attr.flags & QSE_HTRE_ATTR_CHUNKED))
{ {
@ -2087,21 +2094,24 @@ if (qse_htre_getcontentlen(req) > 0)
qse_httpd_discardcontent (httpd, req); qse_httpd_discardcontent (httpd, req);
} }
} }
if (task == QSE_NULL) goto oops; if (task == QSE_NULL) goto oops;
} }
}
else else
{ {
/* contents are all received */ /* contents are all received */
if (mth == QSE_HTTP_CONNECT) if (mth == QSE_HTTP_CONNECT)
{ {
printf ("SWITCHING HTRD TO DUMMY....\n"); printf ("SWITCHING HTRD TO DUMMY.... %s\n", qse_htre_getqpath(req));
/* Switch the http read to a dummy mode so that the subsqeuent /* Switch the http read to a dummy mode so that the subsqeuent
* input is just treaet as connects to the request just completed */ * input is just treaet as connects to the request just completed */
qse_htrd_setoption (client->htrd, qse_htrd_getoption(client->htrd) | QSE_HTRD_DUMMY); qse_htrd_setoption (client->htrd, qse_htrd_getoption(client->htrd) | QSE_HTRD_DUMMY);
if (server_xtn->makersrc (httpd, client, req, &rsrc) <= -1) if (server_xtn->makersrc (httpd, client, req, &rsrc) <= -1)
{ {
printf ("CANOT MAKE RESOURCE.... %s\n", qse_htre_getqpath(req));
/* failed to make a resource. just send the internal server error. /* failed to make a resource. just send the internal server error.
* the makersrc handler can return a negative number to return * the makersrc handler can return a negative number to return
* '500 Internal Server Error'. If it wants to return a specific * '500 Internal Server Error'. If it wants to return a specific
@ -2547,7 +2557,7 @@ static int make_resource (
target->u.proxy.src.type = target->u.proxy.dst.type; target->u.proxy.src.type = target->u.proxy.dst.type;
/* mark that this request is going to be proxied. */ /* mark that this request is going to be proxied. */
/*req->attr.flags |= QSE_HTRE_ATTR_PROXIED;*/ req->attr.flags |= QSE_HTRE_ATTR_PROXIED;
return 0; return 0;
} }

View File

@ -176,7 +176,7 @@ int qse_httpd_setopt (qse_httpd_t* httpd, qse_httpd_opt_t id, const void* value)
return -1; return -1;
} }
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
qse_httpd_ecb_t* qse_httpd_popecb (qse_httpd_t* httpd) qse_httpd_ecb_t* qse_httpd_popecb (qse_httpd_t* httpd)
{ {
@ -191,7 +191,7 @@ void qse_httpd_pushecb (qse_httpd_t* httpd, qse_httpd_ecb_t* ecb)
httpd->ecb = ecb; httpd->ecb = ecb;
} }
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
QSE_INLINE void* qse_httpd_allocmem (qse_httpd_t* httpd, qse_size_t size) QSE_INLINE void* qse_httpd_allocmem (qse_httpd_t* httpd, qse_size_t size)
{ {
@ -249,7 +249,7 @@ qse_mchar_t* qse_httpd_strntombsdup (qse_httpd_t* httpd, const qse_char_t* str,
return mptr; return mptr;
} }
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
static qse_httpd_real_task_t* enqueue_task ( static qse_httpd_real_task_t* enqueue_task (
qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_t* httpd, qse_httpd_client_t* client,
@ -357,7 +357,7 @@ static QSE_INLINE void purge_tasks (
while (dequeue_task (httpd, client) == 0); while (dequeue_task (httpd, client) == 0);
} }
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
static int htrd_peek_request (qse_htrd_t* htrd, qse_htre_t* req) static int htrd_peek_request (qse_htrd_t* htrd, qse_htre_t* req)
{ {
@ -377,7 +377,7 @@ static qse_htrd_recbs_t htrd_recbs =
QSE_STRUCT_FIELD (poke, htrd_poke_request) QSE_STRUCT_FIELD (poke, htrd_poke_request)
}; };
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
static qse_httpd_client_t* new_client ( static qse_httpd_client_t* new_client (
qse_httpd_t* httpd, qse_httpd_client_t* tmpl) qse_httpd_t* httpd, qse_httpd_client_t* tmpl)
@ -580,7 +580,7 @@ qse_printf (QSE_T("MUX ADDHND CLIENT READ %d\n"), client->handle.i);
return 0; return 0;
} }
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
static void deactivate_servers (qse_httpd_t* httpd) static void deactivate_servers (qse_httpd_t* httpd)
{ {
@ -726,7 +726,37 @@ qse_httpd_server_t* qse_httpd_getprevserver (qse_httpd_t* httpd, qse_httpd_serve
return server->prev; return server->prev;
} }
/* --------------------------------------------------- */ /* ----------------------------------------------------------------------- */
#if 0
qse_httpd_dns_t* qse_httpd_attachdns (qse_httpd_t* httpd, qse_httpd_dns_dope_t* dns, qse_size_t xtnsize)
{
qse_httpd_dns_t* dns;
dns = qse_httpd_callocmem (httpd, QSE_SIZEOF(*dns) + xtnsize);
if (dns == QSE_NULL) return QSE_NULL;
dns->type = QSE_HTTPD_SERVER;
/* copy the dns dope */
dns->dope = *dope;
/* and correct some fields in case the dope contains invalid stuffs */
dns->dope.flags &= ~QSE_HTTPD_SERVER_ACTIVE;
/* chain the dns to the tail of the list */
dns->prev = httpd->dns.list.tail;
dns->next = QSE_NULL;
if (httpd->dns.list.tail)
httpd->dns.list.tail->next = dns;
else
httpd->dns.list.head = dns;
httpd->dns.list.tail = dns;
httpd->dns.navail++;
return dns;
}
#endif
/* ----------------------------------------------------------------------- */
static int read_from_client (qse_httpd_t* httpd, qse_httpd_client_t* client) static int read_from_client (qse_httpd_t* httpd, qse_httpd_client_t* client)
{ {
@ -858,143 +888,14 @@ qse_printf (QSE_T("!!!!!FEEDING OK OK OK OK %d from %d\n"), (int)m, (int)client-
return 0; return 0;
} }
static int invoke_client_task ( static int update_mux_for_current_task (qse_httpd_t* httpd, qse_httpd_client_t* client, qse_httpd_task_t* task)
qse_httpd_t* httpd, qse_httpd_client_t* client,
qse_ubi_t handle, int mask)
{ {
qse_httpd_task_t* task;
qse_size_t i; qse_size_t i;
int n, trigger_fired, client_handle_writable;
/* TODO: handle comparison callback ... */
if (handle.i == client->handle.i && (mask & QSE_HTTPD_MUX_READ)) /* TODO: no direct comparision */
{
if (!(client->status & CLIENT_MUTE) &&
read_from_client (httpd, client) <= -1)
{
/* return failure on disconnection also in order to
* purge the client in perform_client_task().
* thus the following line isn't necessary.
*if (httpd->errnum == QSE_HTTPD_EDISCON) return 0;*/
return -1;
}
}
/* this client doesn't have any task */
task = client->task.head;
if (task == QSE_NULL)
{
if (client->status & CLIENT_MUTE)
{
/* handle this delayed client disconnection */
return -1;
}
return 0;
}
trigger_fired = 0;
client_handle_writable = 0;
for (i = 0; i < QSE_COUNTOF(task->trigger); i++)
{
task->trigger[i].mask &= ~(QSE_HTTPD_TASK_TRIGGER_READABLE |
QSE_HTTPD_TASK_TRIGGER_WRITABLE);
if (task->trigger[i].handle.i == handle.i) /* TODO: no direct comparision */
{
if (task->trigger[i].mask & QSE_HTTPD_TASK_TRIGGER_READ)
{
trigger_fired = 1;
task->trigger[i].mask |= QSE_HTTPD_TASK_TRIGGER_READABLE;
}
if (task->trigger[i].mask & QSE_HTTPD_TASK_TRIGGER_WRITE)
{
trigger_fired = 1;
task->trigger[i].mask |= QSE_HTTPD_TASK_TRIGGER_WRITABLE;
if (handle.i == client->handle.i) client_handle_writable = 1; /* TODO: no direct comparison */
}
}
}
if (trigger_fired && !client_handle_writable)
{
/* the task is invoked for triggers.
* check if the client handle is writable */
qse_ntime_t tmout;
tmout.sec = 0;
tmout.nsec = 0;
if (httpd->opt.scb.mux.writable (httpd, client->handle, &tmout) <= 0)
{
/* it is not writable yet. so just skip
* performing the actual task */
return 0;
}
}
n = task->main (httpd, client, task);
if (n <= -1) return -1;
else if (n == 0)
{
int mux_mask;
int mux_status;
/* the current task is over. remove the task
* from the queue. dequeue_task() clears task triggers
* from the mux. so i don't clear them explicitly here */
dequeue_task (httpd, client);
mux_mask = QSE_HTTPD_MUX_READ;
mux_status = CLIENT_HANDLE_READ_IN_MUX;
if (client->task.head)
{
/* there is a pending task. arrange to
* trigger it as if it is just entasked */
mux_mask |= QSE_HTTPD_MUX_WRITE;
mux_status |= CLIENT_HANDLE_WRITE_IN_MUX;
if (client->status & CLIENT_MUTE)
{
mux_mask &= ~QSE_HTTPD_MUX_READ;
mux_status &= ~CLIENT_HANDLE_READ_IN_MUX;
}
}
else
{
if (client->status & CLIENT_MUTE)
{
/* no more task. but this client
* has closed connection previously */
return -1;
}
}
if ((client->status & CLIENT_HANDLE_IN_MUX) !=
(mux_status & CLIENT_HANDLE_IN_MUX))
{
httpd->opt.scb.mux.delhnd (httpd, httpd->mux, client->handle);
client->status &= ~CLIENT_HANDLE_IN_MUX;
if (mux_status)
{
if (httpd->opt.scb.mux.addhnd (
httpd, httpd->mux, client->handle, mux_mask, client) <= -1)
{
return -1;
}
client->status |= mux_status;
}
}
QSE_MEMSET (client->trigger, 0, QSE_SIZEOF(client->trigger));
return 0;
}
else
{
/* the code here is pretty fragile. there is a high chance /* the code here is pretty fragile. there is a high chance
* that something can go wrong if the task handler plays * that something can go wrong if the task handler plays
* with the trigger field in an unexpected manner. * with the trigger field in an unexpected manner.
*/ */
for (i = 0; i < QSE_COUNTOF(task->trigger); i++) for (i = 0; i < QSE_COUNTOF(task->trigger); i++)
{ {
task->trigger[i].mask &= ~(QSE_HTTPD_TASK_TRIGGER_READABLE | task->trigger[i].mask &= ~(QSE_HTTPD_TASK_TRIGGER_READABLE |
@ -1110,6 +1011,151 @@ static int invoke_client_task (
} }
return 0; return 0;
} }
static int update_mux_for_next_task (qse_httpd_t* httpd, qse_httpd_client_t* client)
{
int mux_mask;
int mux_status;
mux_mask = QSE_HTTPD_MUX_READ;
mux_status = CLIENT_HANDLE_READ_IN_MUX;
if (client->task.head)
{
/* there is a pending task. arrange to
* trigger it as if it is just entasked */
mux_mask |= QSE_HTTPD_MUX_WRITE;
mux_status |= CLIENT_HANDLE_WRITE_IN_MUX;
if (client->status & CLIENT_MUTE)
{
mux_mask &= ~QSE_HTTPD_MUX_READ;
mux_status &= ~CLIENT_HANDLE_READ_IN_MUX;
}
}
else
{
if (client->status & CLIENT_MUTE)
{
/* no more task. but this client
* has closed connection previously */
return -1;
}
}
if ((client->status & CLIENT_HANDLE_IN_MUX) !=
(mux_status & CLIENT_HANDLE_IN_MUX))
{
httpd->opt.scb.mux.delhnd (httpd, httpd->mux, client->handle);
client->status &= ~CLIENT_HANDLE_IN_MUX;
if (mux_status)
{
if (httpd->opt.scb.mux.addhnd (
httpd, httpd->mux, client->handle, mux_mask, client) <= -1)
{
return -1;
}
client->status |= mux_status;
}
}
QSE_MEMSET (client->trigger, 0, QSE_SIZEOF(client->trigger));
return 0;
}
static int invoke_client_task (
qse_httpd_t* httpd, qse_httpd_client_t* client,
qse_ubi_t handle, int mask)
{
qse_httpd_task_t* task;
qse_size_t i;
int n, trigger_fired, client_handle_writable;
/* TODO: handle comparison callback ... */
if (handle.i == client->handle.i && (mask & QSE_HTTPD_MUX_READ)) /* TODO: no direct comparision */
{
if (!(client->status & CLIENT_MUTE) &&
read_from_client (httpd, client) <= -1)
{
/* return failure on disconnection also in order to
* purge the client in perform_client_task().
* thus the following line isn't necessary.
*if (httpd->errnum == QSE_HTTPD_EDISCON) return 0;*/
return -1;
}
}
/* this client doesn't have any task */
task = client->task.head;
if (task == QSE_NULL)
{
if (client->status & CLIENT_MUTE)
{
/* handle this delayed client disconnection */
return -1;
}
return 0;
}
trigger_fired = 0;
client_handle_writable = 0;
for (i = 0; i < QSE_COUNTOF(task->trigger); i++)
{
task->trigger[i].mask &= ~(QSE_HTTPD_TASK_TRIGGER_READABLE |
QSE_HTTPD_TASK_TRIGGER_WRITABLE);
if (task->trigger[i].handle.i == handle.i) /* TODO: no direct comparision */
{
if (mask & QSE_HTTPD_MUX_READ)
{
QSE_ASSERT (task->trigger[i].mask & QSE_HTTPD_TASK_TRIGGER_READ);
trigger_fired = 1;
task->trigger[i].mask |= QSE_HTTPD_TASK_TRIGGER_READABLE;
}
if (mask & QSE_HTTPD_MUX_WRITE)
{
QSE_ASSERT (task->trigger[i].mask & QSE_HTTPD_TASK_TRIGGER_WRITE);
trigger_fired = 1;
task->trigger[i].mask |= QSE_HTTPD_TASK_TRIGGER_WRITABLE;
if (handle.i == client->handle.i) client_handle_writable = 1; /* TODO: no direct comparison */
}
}
}
if (trigger_fired && !client_handle_writable)
{
/* the task is invoked for triggers.
* check if the client handle is writable */
qse_ntime_t tmout;
tmout.sec = 0;
tmout.nsec = 0;
if (httpd->opt.scb.mux.writable (httpd, client->handle, &tmout) <= 0)
{
/* it is not writable yet. so just skip
* performing the actual task */
return 0;
}
}
n = task->main (httpd, client, task);
if (n <= -1)
{
/* task error */
return -1;
}
else if (n == 0)
{
/* the current task is over. remove the task
* from the queue. dequeue_task() clears task triggers
* from the mux. so i don't clear them explicitly here */
dequeue_task (httpd, client);
return update_mux_for_next_task (httpd, client);
}
else
{
return update_mux_for_current_task (httpd, client, task);
}
} }
static int perform_client_task ( static int perform_client_task (
@ -1219,7 +1265,7 @@ qse_httpd_task_t* qse_httpd_entask (
client->status &= ~CLIENT_HANDLE_IN_MUX; client->status &= ~CLIENT_HANDLE_IN_MUX;
#if 0 #if 0
qse_printf (QSE_T("MUX ADDHND CLIENT RW(ENTASK) %d\n"), client->handle.i); printf ("MUX ADDHND CLIENT RW(ENTASK) %d\n", client->handle.i);
#endif #endif
if (httpd->opt.scb.mux.addhnd ( if (httpd->opt.scb.mux.addhnd (
httpd, httpd->mux, client->handle, httpd, httpd->mux, client->handle,