Feersum
view release on metacpan or search on metacpan
struct phr_header headers[MAX_HEADERS];
SV* path;
SV* query;
SV* addr;
SV* port;
};
enum feer_respond_state {
RESPOND_NOT_STARTED = 0,
RESPOND_NORMAL = 1,
RESPOND_STREAMING = 2,
RESPOND_SHUTDOWN = 3
};
#define RESPOND_STR(_n,_s) do { \
switch(_n) { \
case RESPOND_NOT_STARTED: _s = "NOT_STARTED(0)"; break; \
case RESPOND_NORMAL: _s = "NORMAL(1)"; break; \
case RESPOND_STREAMING: _s = "STREAMING(2)"; break; \
case RESPOND_SHUTDOWN: _s = "SHUTDOWN(4)"; break; \
} \
} while (0)
enum feer_receive_state {
RECEIVE_WAIT = 0,
RECEIVE_HEADERS = 1,
RECEIVE_BODY = 2,
RECEIVE_STREAMING = 3,
RECEIVE_SHUTDOWN = 4
};
#define RECEIVE_STR(_n,_s) do { \
switch(_n) { \
case RECEIVE_WAIT: _s = "WAIT(0)"; break; \
case RECEIVE_HEADERS: _s = "HEADERS(1)"; break; \
case RECEIVE_BODY: _s = "BODY(2)"; break; \
case RECEIVE_STREAMING: _s = "STREAMING(3)"; break; \
case RECEIVE_SHUTDOWN: _s = "SHUTDOWN(4)"; break; \
} \
} while (0)
struct feer_conn {
SV *self;
int fd;
struct sockaddr *sa;
struct ev_io read_ev_io;
struct ev_io write_ev_io;
struct ev_timer read_ev_timer;
SV *rbuf;
struct rinq *wbuf_rinq;
SV *poll_write_cb;
SV *ext_guard;
struct feer_req *req;
ssize_t expected_cl;
ssize_t received_cl;
enum feer_respond_state responding;
enum feer_receive_state receiving;
bool is_keepalive;
int reqs;
unsigned int in_callback;
unsigned int is_http11:1;
unsigned int poll_write_cb_is_io_handle:1;
unsigned int auto_cl:1;
ssize_t pipelined;
};
enum feer_header_norm_style {
HEADER_NORM_SKIP = 0,
HEADER_NORM_UPCASE_DASH = 1,
HEADER_NORM_LOCASE_DASH = 2,
HEADER_NORM_UPCASE = 3,
HEADER_NORM_LOCASE = 4
};
typedef struct feer_conn feer_conn_handle; // for typemap
#define dCONN struct feer_conn *c = (struct feer_conn *)w->data
#define IsArrayRef(_x) (SvROK(_x) && SvTYPE(SvRV(_x)) == SVt_PVAV)
#define IsCodeRef(_x) (SvROK(_x) && SvTYPE(SvRV(_x)) == SVt_PVCV)
static SV* feersum_env_method(pTHX_ struct feer_req *r);
static SV* feersum_env_uri(pTHX_ struct feer_req *r);
static SV* feersum_env_protocol(pTHX_ struct feer_req *r);
static void feersum_set_path_and_query(pTHX_ struct feer_req *r);
static void feersum_set_remote_info(pTHX_ struct feer_req *r, struct sockaddr *sa);
static HV* feersum_env(pTHX_ struct feer_conn *c);
static SV* feersum_env_path(pTHX_ struct feer_req *r);
static SV* feersum_env_query(pTHX_ struct feer_req *r);
static HV* feersum_env_headers(pTHX_ struct feer_req *r, int norm);
static SV* feersum_env_header(pTHX_ struct feer_req *r, SV* name);
static SV* feersum_env_addr(pTHX_ struct feer_conn *c);
static SV* feersum_env_port(pTHX_ struct feer_conn *c);
static ssize_t feersum_env_content_length(pTHX_ struct feer_conn *c);
static SV* feersum_env_io(pTHX_ struct feer_conn *c);
static void feersum_start_response
(pTHX_ struct feer_conn *c, SV *message, AV *headers, int streaming);
static size_t feersum_write_whole_body (pTHX_ struct feer_conn *c, SV *body);
static void feersum_handle_psgi_response(
pTHX_ struct feer_conn *c, SV *ret, bool can_recurse);
static bool feersum_set_keepalive (pTHX_ struct feer_conn *c, bool is_keepalive);
static int feersum_close_handle(pTHX_ struct feer_conn *c, bool is_writer);
static SV* feersum_conn_guard(pTHX_ struct feer_conn *c, SV *guard);
static void start_read_watcher(struct feer_conn *c);
static void stop_read_watcher(struct feer_conn *c);
static void restart_read_timer(struct feer_conn *c);
static void stop_read_timer(struct feer_conn *c);
static void start_write_watcher(struct feer_conn *c);
static void stop_write_watcher(struct feer_conn *c);
static void try_conn_write(EV_P_ struct ev_io *w, int revents);
static void try_conn_read(EV_P_ struct ev_io *w, int revents);
static void conn_read_timeout(EV_P_ struct ev_timer *w, int revents);
static bool process_request_headers(struct feer_conn *c, int body_offset);
static void sched_request_callback(struct feer_conn *c);
static void call_died (pTHX_ struct feer_conn *c, const char *cb_type);
static void call_request_callback(struct feer_conn *c);
static void call_poll_callback (struct feer_conn *c, bool is_write);
static void pump_io_handle (struct feer_conn *c, SV *io);
static void conn_write_ready (struct feer_conn *c);
static void respond_with_server_error(struct feer_conn *c, const char *msg, STRLEN msg_len, int code);
static void update_wbuf_placeholder(struct feer_conn *c, SV *sv, struct iovec *iov);
static STRLEN add_sv_to_wbuf (struct feer_conn *c, SV *sv);
static STRLEN add_const_to_wbuf (struct feer_conn *c, const char *str, size_t str_len);
#define add_crlf_to_wbuf(c) add_const_to_wbuf(c,CRLF,2)
static void finish_wbuf (struct feer_conn *c);
static void add_chunk_sv_to_wbuf (struct feer_conn *c, SV *sv);
static void add_placeholder_to_wbuf (struct feer_conn *c, SV **sv, struct iovec **iov_ref);
static void uri_decode_sv (SV *sv);
static bool str_eq(const char *a, int a_len, const char *b, int b_len);
static bool str_case_eq(const char *a, int a_len, const char *b, int b_len);
static SV* fetch_av_normal (pTHX_ AV *av, I32 i);
static const char *http_code_to_msg (int code);
static int prep_socket (int fd, int is_tcp);
static HV *feer_stash, *feer_conn_stash;
static HV *feer_conn_reader_stash = NULL, *feer_conn_writer_stash = NULL;
static MGVTBL psgix_io_vtbl;
static SV *request_cb_cv = NULL;
static bool request_cb_is_psgi = 0;
static SV *shutdown_cb_cv = NULL;
static bool shutting_down = 0;
static int active_conns = 0;
static double read_timeout = READ_TIMEOUT;
static unsigned int max_connection_reqs = 0;
static SV *feer_server_name = NULL;
static SV *feer_server_port = NULL;
static bool is_tcp = 1;
static bool is_keepalive = KEEPALIVE_CONNECTION;
static ev_io accept_w;
static ev_prepare ep;
static ev_check ec;
struct ev_idle ei;
static struct rinq *request_ready_rinq = NULL;
static AV *psgi_ver;
static SV *psgi_serv10, *psgi_serv11, *crlf_sv;
// TODO: make this thread-local if and when there are multiple C threads:
struct ev_loop *feersum_ev_loop = NULL;
static HV *feersum_tmpl_env = NULL;
#define DATE_HEADER_LENGTH 37 // "Date: Thu, 01 Jan 1970 00:00:00 GMT\015\012"
static const char *const DAYS[] = {"Sun", "Mon", "Tue", "Wed", "Thu", "Fri", "Sat"};
static const char *const MONTHS[] = {"Jan", "Feb", "Mar", "Apr", "May", "Jun",
"Jul", "Aug", "Sep", "Oct", "Nov", "Dec"};
static char DATE_BUF[DATE_HEADER_LENGTH+1] = "Date: The, 01 Jan 1970 00:00:00 GMT\015\012";
static time_t LAST_GENERATED_TIME = 0;
static INLINE_UNLESS_DEBUG void uint_to_str(unsigned int value, char *str) {
str[0] = (value / 10) + '0';
str[1] = (value % 10) + '0';
}
static INLINE_UNLESS_DEBUG void uint_to_str_4digits(unsigned int value, char *str) {
str[0] = (value / 1000) + '0';
str[1] = (value / 100) % 10 + '0';
str[2] = (value / 10) % 10 + '0';
str[3] = value % 10 + '0';
}
INLINE_UNLESS_DEBUG
static void generate_date_header(void) {
time_t now = time(NULL);
if (now == LAST_GENERATED_TIME) return;
LAST_GENERATED_TIME = now;
struct tm *tm = gmtime(&now);
const char *day = DAYS[tm->tm_wday];
DATE_BUF[6] = day[0];
DATE_BUF[7] = day[1];
DATE_BUF[8] = day[2];
uint_to_str(tm->tm_mday, DATE_BUF + 11);
const char *month = MONTHS[tm->tm_mon];
DATE_BUF[14] = month[0];
DATE_BUF[15] = month[1];
DATE_BUF[16] = month[2];
uint_to_str_4digits(tm->tm_year + 1900, DATE_BUF + 18);
uint_to_str(tm->tm_hour, DATE_BUF + 23);
uint_to_str(tm->tm_min, DATE_BUF + 26);
uint_to_str(tm->tm_sec, DATE_BUF + 29);
int flags;
// make it non-blocking
flags = O_NONBLOCK;
if (unlikely(fcntl(fd, F_SETFL, flags) < 0))
return -1;
flags = 1;
#endif
if (likely(is_tcp)) {
// flush writes immediately
if (unlikely(setsockopt(fd, SOL_TCP, TCP_NODELAY, &flags, sizeof(int))))
return -1;
}
// handle URG data inline
if (unlikely(setsockopt(fd, SOL_SOCKET, SO_OOBINLINE, &flags, sizeof(int))))
return -1;
// disable lingering
struct linger linger = { .l_onoff = 0, .l_linger = 0 };
if (unlikely(setsockopt(fd, SOL_SOCKET, SO_LINGER, &linger, sizeof(linger))))
return -1;
return 0;
}
INLINE_UNLESS_DEBUG static void
safe_close_conn(struct feer_conn *c, const char *where)
{
if (unlikely(c->fd < 0))
return;
// make it blocking
fcntl(c->fd, F_SETFL, 0);
if (unlikely(close(c->fd)))
perror(where);
c->fd = -1;
}
static struct feer_conn *
new_feer_conn (EV_P_ int conn_fd, struct sockaddr *sa)
{
SV *self = newSV(0);
SvUPGRADE(self, SVt_PVMG); // ensures sv_bless doesn't reallocate
SvGROW(self, sizeof(struct feer_conn));
SvPOK_only(self);
SvIOK_on(self);
SvIV_set(self,conn_fd);
struct feer_conn *c = (struct feer_conn *)SvPVX(self);
Zero(c, 1, struct feer_conn);
c->self = self;
c->fd = conn_fd;
c->sa = sa;
c->responding = RESPOND_NOT_STARTED;
c->receiving = RECEIVE_HEADERS;
c->is_keepalive = 0;
c->reqs = 0;
c->pipelined = 0;
ev_io_init(&c->read_ev_io, try_conn_read, conn_fd, EV_READ);
c->read_ev_io.data = (void *)c;
ev_init(&c->read_ev_timer, conn_read_timeout);
c->read_ev_timer.data = (void *)c;
trace3("made conn fd=%d self=%p, c=%p, cur=%"Sz_uf", len=%"Sz_uf"\n",
c->fd, self, c, (Sz)SvCUR(self), (Sz)SvLEN(self));
SV *rv = newRV_inc(c->self);
sv_bless(rv, feer_conn_stash); // so DESTROY can get called on read errors
SvREFCNT_dec(rv);
SvREADONLY_on(self); // turn off later for blessing
active_conns++;
return c;
}
// for use in the typemap:
INLINE_UNLESS_DEBUG
static struct feer_conn *
sv_2feer_conn (SV *rv)
{
if (unlikely(!sv_isa(rv,"Feersum::Connection")))
croak("object is not of type Feersum::Connection");
return (struct feer_conn *)SvPVX(SvRV(rv));
}
INLINE_UNLESS_DEBUG
static SV*
feer_conn_2sv (struct feer_conn *c)
{
return newRV_inc(c->self);
}
static feer_conn_handle *
sv_2feer_conn_handle (SV *rv, bool can_croak)
{
trace3("sv 2 conn_handle\n");
if (unlikely(!SvROK(rv))) croak("Expected a reference");
// do not allow subclassing
SV *sv = SvRV(rv);
if (likely(
sv_isobject(rv) &&
(SvSTASH(sv) == feer_conn_writer_stash ||
SvSTASH(sv) == feer_conn_reader_stash)
)) {
UV uv = SvUV(sv);
if (uv == 0) {
if (can_croak) croak("Operation not allowed: Handle is closed.");
return NULL;
}
return INT2PTR(feer_conn_handle*,uv);
}
if (can_croak)
croak("Expected a Feersum::Connection::Writer or ::Reader object");
w->fd, v->iov_base, (Sz)v->iov_len);
v->iov_base += wrote;
v->iov_len -= wrote;
// don't consume any more:
consume = 0;
}
else {
trace3("consume vector %d base=%p len=%"Sz_uf" sv=%p\n",
w->fd, v->iov_base, (Sz)v->iov_len, m->sv[i]);
wrote -= v->iov_len;
m->offset++;
if (m->sv[i]) {
SvREFCNT_dec(m->sv[i]);
m->sv[i] = NULL;
}
}
}
if (likely(m->offset >= m->count)) {
trace2("all done with iomatrix %d state=%d\n",w->fd,c->responding);
rinq_shift(&c->wbuf_rinq);
Safefree(m);
if (!c->wbuf_rinq)
goto try_write_finished;
trace2("write again immediately %d state=%d\n",w->fd,c->responding);
goto try_write_again_immediately;
}
// else, fallthrough:
trace2("write fallthrough %d state=%d\n",w->fd,c->responding);
try_write_again:
trace("write again %d state=%d\n",w->fd,c->responding);
start_write_watcher(c);
goto try_write_cleanup;
try_write_finished:
// should always be responding, but just in case
switch(c->responding) {
case RESPOND_NOT_STARTED:
// the write watcher shouldn't ever get called before starting to
// respond. Shut it down if it does.
trace("unexpected try_write when response not started %d\n",c->fd);
goto try_write_shutdown;
case RESPOND_NORMAL:
goto try_write_shutdown;
case RESPOND_STREAMING:
if (c->poll_write_cb) goto try_write_again;
else goto try_write_paused;
case RESPOND_SHUTDOWN:
goto try_write_shutdown;
default:
goto try_write_cleanup;
}
try_write_paused:
trace3("write PAUSED %d, refcnt=%d, state=%d\n", c->fd, SvREFCNT(c->self), c->responding);
stop_write_watcher(c);
goto try_write_cleanup;
try_write_shutdown:
if (likely(c->is_keepalive)) {
trace3("write SHUTDOWN, but KEEP %d, refcnt=%d, state=%d\n", c->fd, SvREFCNT(c->self), c->responding);
stop_write_watcher(c);
change_responding_state(c, RESPOND_NOT_STARTED);
change_receiving_state(c, RECEIVE_WAIT);
if (likely(c->req)) {
if (c->req->buf) SvREFCNT_dec(c->req->buf);
if (likely(c->req->path)) SvREFCNT_dec(c->req->path);
if (likely(c->req->query)) SvREFCNT_dec(c->req->query);
if (likely(c->req->addr)) SvREFCNT_dec(c->req->addr);
if (likely(c->req->port)) SvREFCNT_dec(c->req->port);
Safefree(c->req);
}
c->req = NULL;
ssize_t pipelined = 0;
if (c->rbuf) { pipelined = SvCUR(c->rbuf); }
if (unlikely(pipelined > 0 && c->is_http11)) {
trace3("connections has pipelined data on %d\n", c->fd);
c->pipelined = pipelined;
try_conn_read(EV_A, &c->read_ev_io, 0);
} else {
c->pipelined = 0;
start_read_watcher(c);
restart_read_timer(c);
}
trace3("connections active on %d\n", c->fd);
} else {
trace3("write SHUTDOWN %d, refcnt=%d, state=%d\n", c->fd, SvREFCNT(c->self), c->responding);
stop_write_watcher(c);
change_responding_state(c, RESPOND_SHUTDOWN);
safe_close_conn(c, "close at write shutdown");
}
try_write_cleanup:
SvREFCNT_dec(c->self);
return;
}
static int
try_parse_http(struct feer_conn *c, size_t last_read)
{
struct feer_req *req = c->req;
if (likely(!req)) {
Newxz(req,1,struct feer_req);
c->req = req;
}
// GH#12 - incremental parsing sets num_headers to 0 each time; force it
// back on every invocation
req->num_headers = MAX_HEADERS;
return phr_parse_request(SvPVX(c->rbuf), SvCUR(c->rbuf),
&req->method, &req->method_len,
&req->uri, &req->uri_len, &req->minor_version,
req->headers, &req->num_headers,
(SvCUR(c->rbuf)-last_read));
}
static void
try_conn_read(EV_P_ ev_io *w, int revents)
{
dCONN;
SvREFCNT_inc_void_NN(c->self);
ssize_t got_n = 0;
if (unlikely(c->pipelined)) goto pipelined;
// if it's marked readable EV suggests we simply try read it. Otherwise it
// is stopped and we should ditch this connection.
if (unlikely(revents & EV_ERROR && !(revents & EV_READ))) {
trace("EV error on read, fd=%d revents=0x%08x\n", w->fd, revents);
goto try_read_error;
}
if (unlikely(c->receiving == RECEIVE_SHUTDOWN))
goto dont_read_again;
trace("try read %d\n",w->fd);
if (unlikely(!c->rbuf)) { // likely = optimize for keepalive requests
trace("init rbuf for %d\n",w->fd);
c->rbuf = newSV(READ_INIT_FACTOR*READ_BUFSZ + 1);
SvPOK_on(c->rbuf);
}
ssize_t space_free = SvLEN(c->rbuf) - SvCUR(c->rbuf);
if (unlikely(space_free < READ_BUFSZ)) { // unlikely = optimize for small
size_t new_len = SvLEN(c->rbuf) + READ_GROW_FACTOR*READ_BUFSZ;
trace("moar memory %d: %"Sz_uf" to %"Sz_uf"\n",
w->fd, (Sz)SvLEN(c->rbuf), (Sz)new_len);
SvGROW(c->rbuf, new_len);
space_free += READ_GROW_FACTOR*READ_BUFSZ;
}
char *cur = SvPVX(c->rbuf) + SvCUR(c->rbuf);
got_n = read(w->fd, cur, space_free);
if (unlikely(got_n <= 0)) {
if (unlikely(got_n == 0)) {
trace("EOF before complete request: %d\n",w->fd,SvCUR(c->rbuf));
goto try_read_error;
}
if (likely(errno == EAGAIN || errno == EINTR))
goto try_read_again;
perror("try_conn_read error");
goto try_read_error;
}
trace("read %d %"Ssz_df"\n", w->fd, (Ssz)got_n);
SvCUR(c->rbuf) += got_n;
goto try_parse;
pipelined:
got_n = c->pipelined;
c->pipelined = 0;
try_parse:
// likely = optimize for small requests
if (likely(c->receiving <= RECEIVE_HEADERS)) {
int ret = try_parse_http(c, (size_t)got_n);
if (ret == -1) goto try_read_bad;
#ifdef TCP_DEFER_ACCEPT
if (ret == -2) goto try_read_again_reset_timer;
#else
if (ret == -2) {
if (is_tcp) goto try_read_again;
else goto try_read_again_reset_timer;
}
#endif
if (process_request_headers(c, ret))
goto try_read_again_reset_timer;
else
goto dont_read_again;
}
else if (likely(c->receiving == RECEIVE_BODY)) {
c->received_cl += got_n;
if (c->received_cl < c->expected_cl)
goto try_read_again_reset_timer;
// body is complete
sched_request_callback(c);
goto dont_read_again;
}
else {
trouble("unknown read state %d %d", w->fd, c->receiving);
}
// fallthrough:
try_read_error:
trace("READ ERROR %d, refcnt=%d\n", w->fd, SvREFCNT(c->self));
change_receiving_state(c, RECEIVE_SHUTDOWN);
change_responding_state(c, RESPOND_SHUTDOWN);
stop_read_watcher(c);
stop_read_timer(c);
stop_write_watcher(c);
goto try_read_cleanup;
try_read_bad:
trace("bad request %d\n", w->fd);
respond_with_server_error(c, "Malformed request.\n", 0, 400);
// TODO: when keep-alive, close conn instead of fallthrough here.
// fallthrough:
dont_read_again:
trace("done reading %d\n", w->fd);
change_receiving_state(c, RECEIVE_SHUTDOWN);
stop_read_watcher(c);
stop_read_timer(c);
goto try_read_cleanup;
try_read_again_reset_timer:
trace("(reset read timer) %d\n", w->fd);
restart_read_timer(c);
// fallthrough:
try_read_again:
trace("read again %d\n", w->fd);
start_read_watcher(c);
try_read_cleanup:
SvREFCNT_dec(c->self);
}
static void
conn_read_timeout (EV_P_ ev_timer *w, int revents)
{
dCONN;
SvREFCNT_inc_void_NN(c->self);
if (unlikely(!(revents & EV_TIMER) || c->receiving == RECEIVE_SHUTDOWN)) {
// if there's no EV_TIMER then EV has stopped it on an error
if (revents & EV_ERROR)
trouble("EV error on read timer, fd=%d revents=0x%08x\n",
c->fd,revents);
goto read_timeout_cleanup;
}
trace("read timeout %d\n", c->fd);
if (likely(c->responding == RESPOND_NOT_STARTED) && c->receiving >= RECEIVE_HEADERS) {
const char *msg;
if (c->receiving == RECEIVE_HEADERS) {
msg = "Headers took too long.";
}
else {
msg = "Timeout reading body.";
}
respond_with_server_error(c, msg, 0, 408);
} else {
trace("read timeout in keepalive conn: %d\n", c->fd);
stop_write_watcher(c);
stop_read_watcher(c);
stop_read_timer(c);
safe_close_conn(c, "close at read timeout");
change_responding_state(c, RESPOND_SHUTDOWN);
}
read_timeout_cleanup:
stop_read_watcher(c);
stop_read_timer(c);
SvREFCNT_dec(c->self);
}
static void
accept_cb (EV_P_ ev_io *w, int revents)
{
struct sockaddr_storage sa_buf;
socklen_t sa_len;
if (unlikely(shutting_down)) {
// shouldn't get called, but be defensive
ev_io_stop(EV_A, w);
close(w->fd);
return;
}
if (unlikely(revents & EV_ERROR)) {
trouble("EV error in accept_cb, fd=%d, revents=0x%08x\n",w->fd,revents);
ev_break(EV_A, EVBREAK_ALL);
return;
}
trace2("accept! revents=0x%08x\n", revents);
while (1) {
sa_len = sizeof(struct sockaddr_storage);
errno = 0;
#ifdef HAS_ACCEPT4
int fd = accept4(w->fd, (struct sockaddr *)&sa_buf, &sa_len, SOCK_CLOEXEC|SOCK_NONBLOCK);
#else
int fd = accept(w->fd, (struct sockaddr *)&sa_buf, &sa_len);
#endif
trace("accepted fd=%d, errno=%d\n", fd, errno);
if (fd == -1) break;
assert(sa_len <= sizeof(struct sockaddr_storage));
if (unlikely(prep_socket(fd, is_tcp))) {
perror("prep_socket");
trouble("prep_socket failed for %d\n", fd);
close(fd);
continue;
}
struct sockaddr *sa = (struct sockaddr *)malloc(sa_len);
memcpy(sa,&sa_buf,(size_t)sa_len);
struct feer_conn *c = new_feer_conn(EV_A,fd,sa);
#ifdef TCP_DEFER_ACCEPT
try_conn_read(EV_A, &c->read_ev_io, EV_READ);
assert(SvREFCNT(c->self) <= 3);
#else
if (is_tcp) {
start_read_watcher(c);
restart_read_timer(c);
assert(SvREFCNT(c->self) == 3);
} else {
try_conn_read(EV_A, &c->read_ev_io, EV_READ);
assert(SvREFCNT(c->self) <= 3);
}
#endif
SvREFCNT_dec(c->self);
}
}
static void
sched_request_callback (struct feer_conn *c)
{
trace("sched req callback: %d c=%p, head=%p\n", c->fd, c, request_ready_rinq);
rinq_push(&request_ready_rinq, c);
SvREFCNT_inc_void_NN(c->self); // for the rinq
if (!ev_is_active(&ei)) {
ev_idle_start(feersum_ev_loop, &ei);
}
}
// the unlikely/likely annotations here are trying to optimize for GET first
// and POST second. Other entity-body requests are third in line.
static bool
process_request_headers (struct feer_conn *c, int body_offset)
{
int err_code;
const char *err;
struct feer_req *req = c->req;
trace("processing headers %d minor_version=%d\n",c->fd,req->minor_version);
bool body_is_required;
bool next_req_follows = 0;
bool got_content_length = 0;
c->is_http11 = (req->minor_version == 1);
c->is_keepalive = is_keepalive && c->is_http11;
c->reqs++;
change_receiving_state(c, RECEIVE_BODY);
if (likely(str_eq("GET", 3, req->method, req->method_len))) {
// Not supposed to have a body. Additional bytes are either a
// mistake, a websocket negotiation or pipelined requests under
// HTTP/1.1
next_req_follows = 1;
}
else if (likely(str_eq("OPTIONS", 7, req->method, req->method_len))) {
next_req_follows = 1;
}
else if (likely(str_eq("POST", 4, req->method, req->method_len))) {
body_is_required = 1;
}
else if (str_eq("PUT", 3, req->method, req->method_len)) {
body_is_required = 1;
}
else if (str_eq("HEAD", 4, req->method, req->method_len) ||
str_eq("DELETE", 6, req->method, req->method_len))
{
next_req_follows = 1;
}
else {
err = "Feersum doesn't support that method yet\n";
err_code = 405;
goto got_bad_request;
}
#if DEBUG >= 2
if (next_req_follows)
trace2("next req follows fd=%d, boff=%d\n",c->fd,body_offset);
if (body_is_required)
trace2("body is required fd=%d, boff=%d\n",c->fd,body_offset);
#endif
// a body or follow-on data potentially follows the headers. Let feer_req
// retain its pointers into rbuf and make a new scalar for more body data.
STRLEN from_len;
char *from = SvPV(c->rbuf,from_len);
from += body_offset;
int need = from_len - body_offset;
int new_alloc = (need > READ_INIT_FACTOR*READ_BUFSZ)
? need : READ_INIT_FACTOR*READ_BUFSZ-1;
trace("new rbuf for body %d need=%d alloc=%d\n",c->fd, need, new_alloc);
SV *new_rbuf = newSVpvn(need ? from : "", need);
req->buf = c->rbuf;
c->rbuf = new_rbuf;
SvCUR_set(req->buf, body_offset);
// determine how much we need to read
int i;
UV expected = 0;
for (i=0; i < req->num_headers; i++) {
struct phr_header *hdr = &req->headers[i];
if (!hdr->name) continue;
// XXX: ignore multiple C-L headers?
if (unlikely(
str_case_eq("content-length", 14, hdr->name, hdr->name_len)))
{
int g = grok_number(hdr->value, hdr->value_len, &expected);
if (likely(g == IS_NUMBER_IN_UV)) {
if (unlikely(expected > MAX_BODY_LEN)) {
err_code = 413;
err = "Content length exceeds maximum\n";
goto got_bad_request;
}
else
got_content_length = 1;
}
else {
err_code = 400;
err = "invalid content-length\n";
goto got_bad_request;
}
}
else if (
unlikely(str_case_eq("connection", 10, hdr->name, hdr->name_len)))
{
if (likely(c->is_http11)
&& likely(c->is_keepalive)
&& likely(str_case_eq("close", 5, hdr->value, hdr->value_len)))
{
c->is_keepalive = 0;
trace("setting conn %d to close after response\n", c->fd);
}
else if (
likely(!c->is_http11)
&& likely(is_keepalive)
&& str_case_eq("keep-alive", 10, hdr->value, hdr->value_len))
{
c->is_keepalive = 1;
trace("setting conn %d to keep after response\n", c->fd);
}
}
// TODO: support "Transfer-Encoding: chunked" bodies
}
if (max_connection_reqs > 0 && c->reqs >= max_connection_reqs) {
c->is_keepalive = 0;
trace("reached max requests per connection (%d), will close after response\n", max_connection_reqs);
}
if (likely(next_req_follows)) goto got_it_all; // optimize for GET
else if (likely(got_content_length)) goto got_cl;
if (body_is_required) {
// Go the nginx route...
err_code = 411;
err = "Content-Length required\n";
}
else {
// XXX TODO support requests that don't require a body
err_code = 418;
err = "Feersum doesn't know how to handle optional-body requests yet\n";
}
got_bad_request:
respond_with_server_error(c, err, 0, err_code);
return 0;
got_cl:
c->expected_cl = (ssize_t)expected;
c->received_cl = SvCUR(c->rbuf);
trace("expecting body %d size=%"Ssz_df" have=%"Ssz_df"\n",
c->fd, (Ssz)c->expected_cl, (Ssz)c->received_cl);
SvGROW(c->rbuf, c->expected_cl + 1);
// don't have enough bytes to schedule immediately?
// unlikely = optimize for short requests
if (unlikely(c->expected_cl && c->received_cl < c->expected_cl)) {
// TODO: schedule the callback immediately and support a non-blocking
// ->read method.
// sched_request_callback(c);
// change_receiving_state(c, RECEIVE_STREAM);
return 1;
}
// fallthrough: have enough bytes
got_it_all:
sched_request_callback(c);
return 0;
}
static void
conn_write_ready (struct feer_conn *c)
{
if (c->in_callback) return; // defer until out of callback
if (c->write_ev_io.data == NULL) {
ev_io_init(&c->write_ev_io, try_conn_write, c->fd, EV_WRITE);
c->write_ev_io.data = (void *)c;
}
#if AUTOCORK_WRITES
start_write_watcher(c);
#else
// attempt a non-blocking write immediately if we're not already
// waiting for writability
try_conn_write(feersum_ev_loop, &c->write_ev_io, EV_WRITE);
#endif
}
static void
respond_with_server_error (struct feer_conn *c, const char *msg, STRLEN msg_len, int err_code)
{
SV *tmp;
if (unlikely(c->responding != RESPOND_NOT_STARTED)) {
trouble("Tried to send server error but already responding!");
return;
}
if (!msg_len) msg_len = strlen(msg);
assert(msg_len < INT_MAX);
tmp = newSVpvf("HTTP/1.%d %d %s" CRLF
"Content-Type: text/plain" CRLF
"Connection: close" CRLF
"Cache-Control: no-cache, no-store" CRLF
"Content-Length: %"Ssz_df"" CRLFx2
"%.*s",
c->is_http11 ? 1 : 0,
err_code, http_code_to_msg(err_code),
(Ssz)msg_len,
(int)msg_len, msg);
add_sv_to_wbuf(c, sv_2mortal(tmp));
stop_read_watcher(c);
stop_read_timer(c);
change_responding_state(c, RESPOND_SHUTDOWN);
change_receiving_state(c, RECEIVE_SHUTDOWN);
if (c->is_keepalive) c->is_keepalive = 0;
conn_write_ready(c);
}
INLINE_UNLESS_DEBUG bool
str_eq(const char *a, int a_len, const char *b, int b_len)
{
if (a_len != b_len) return 0;
if (a == b) return 1;
int i;
for (i=0; i<a_len && i<b_len; i++) {
if (a[i] != b[i]) return 0;
}
return 1;
}
/*
* Compares two strings, assumes that the first string is already lower-cased
*/
INLINE_UNLESS_DEBUG bool
str_case_eq(const char *a, int a_len, const char *b, int b_len)
{
if (a_len != b_len) return 0;
if (a == b) return 1;
int i;
for (i=0; i<a_len && i<b_len; i++) {
if (a[i] != tolower(b[i])) return 0;
}
return 1;
}
INLINE_UNLESS_DEBUG int
hex_decode(const char ch)
{
if (likely('0' <= ch && ch <= '9'))
return ch - '0';
else if ('A' <= ch && ch <= 'F')
return ch - 'A' + 10;
else if ('a' <= ch && ch <= 'f')
return ch - 'a' + 10;
return -1;
}
static void
uri_decode_sv (SV *sv)
{
STRLEN len;
char *ptr, *end, *decoded;
ptr = SvPV(sv, len);
end = SvEND(sv);
// quickly scan for % so we can ignore decoding that portion of the string
while (ptr < end) {
if (unlikely(*ptr == '%')) goto needs_decode;
ptr++;
}
return;
needs_decode:
if (SvIOK(message))
code = SvIV(message);
else if (SvUOK(message))
code = SvUV(message);
else {
const int numtype = grok_number(SvPVX_const(message),3,&code);
if (unlikely(numtype != IS_NUMBER_IN_UV))
code = 0;
}
trace2("starting response fd=%d code=%"UVuf"\n",c->fd,code);
if (unlikely(!code))
croak("first parameter is not a number or doesn't start with digits");
// for PSGI it's always just an IV so optimize for that
if (likely(!SvPOK(message) || SvCUR(message) == 3)) {
ptr = http_code_to_msg(code);
message = sv_2mortal(newSVpvf("%"UVuf" %s",code,ptr));
}
// don't generate or strip Content-Length headers for 304 or 1xx
c->auto_cl = (code == 304 || code == 204 || (100 <= code && code <= 199)) ? 0 : 1;
add_const_to_wbuf(c, c->is_http11 ? "HTTP/1.1 " : "HTTP/1.0 ", 9);
add_sv_to_wbuf(c, message);
add_crlf_to_wbuf(c);
for (i=0; i<avl; i+= 2) {
SV **hdr = av_fetch(headers, i, 0);
if (unlikely(!hdr || !SvOK(*hdr))) {
trace("skipping undef header key");
continue;
}
SV **val = av_fetch(headers, i+1, 0);
if (unlikely(!val || !SvOK(*val))) {
trace("skipping undef header value");
continue;
}
STRLEN hlen;
const char *hp = SvPV(*hdr, hlen);
if (likely(c->auto_cl) &&
unlikely(str_case_eq("content-length",14,hp,hlen)))
{
trace("ignoring content-length header in the response\n");
continue;
}
add_sv_to_wbuf(c, *hdr);
add_const_to_wbuf(c, ": ", 2);
add_sv_to_wbuf(c, *val);
add_crlf_to_wbuf(c);
}
if (likely(c->is_http11)) {
#ifdef DATE_HEADER
generate_date_header();
add_const_to_wbuf(c, DATE_BUF, DATE_HEADER_LENGTH);
#endif
if (unlikely(!c->is_keepalive))
add_const_to_wbuf(c, "Connection: close" CRLF, 19);
} else if (unlikely(c->is_keepalive) && !streaming)
add_const_to_wbuf(c, "Connection: keep-alive" CRLF, 24);
if (streaming) {
if (c->is_http11)
add_const_to_wbuf(c, "Transfer-Encoding: chunked" CRLFx2, 30);
else {
add_crlf_to_wbuf(c);
// cant do keep-alive for streaming http/1.0 since client completes read on close
if (unlikely(c->is_keepalive)) c->is_keepalive = 0;
}
}
conn_write_ready(c);
}
static size_t
feersum_write_whole_body (pTHX_ struct feer_conn *c, SV *body)
{
size_t RETVAL;
int i;
bool body_is_string = 0;
STRLEN cur;
if (c->responding != RESPOND_NORMAL)
croak("can't use write_whole_body when in streaming mode");
if (!SvOK(body)) {
body = sv_2mortal(newSVpvs(""));
body_is_string = 1;
}
else if (SvROK(body)) {
SV *refd = SvRV(body);
if (SvOK(refd) && !SvROK(refd)) {
body = refd;
body_is_string = 1;
}
else if (SvTYPE(refd) != SVt_PVAV) {
croak("body must be a scalar, scalar reference or array reference");
}
}
else {
body_is_string = 1;
}
SV *cl_sv; // content-length future
struct iovec *cl_iov;
if (likely(c->auto_cl))
add_placeholder_to_wbuf(c, &cl_sv, &cl_iov);
else
add_crlf_to_wbuf(c);
if (body_is_string) {
cur = add_sv_to_wbuf(c,body);
RETVAL = cur;
}
else {
AV *abody = (AV*)SvRV(body);
I32 amax = av_len(abody);
RETVAL = 0;
for (i=0; i<=amax; i++) {
SV *sv = fetch_av_normal(aTHX_ abody, i);
if (unlikely(!sv)) continue;
cur = add_sv_to_wbuf(c,sv);
trace("body part i=%d sv=%p cur=%"Sz_uf"\n", i, sv, (Sz)cur);
RETVAL += cur;
}
}
if (likely(c->auto_cl)) {
trouble("Couldn't close body IO handle: %-p",ERRSV);
}
SvREFCNT_dec(c->poll_write_cb);
c->poll_write_cb = NULL;
finish_wbuf(c);
change_responding_state(c, RESPOND_SHUTDOWN);
goto done_pump_io;
}
if (c->is_http11)
add_chunk_sv_to_wbuf(c, ret);
else
add_sv_to_wbuf(c, ret);
done_pump_io:
trace("leaving pump io handle %d\n", c->fd);
PUTBACK;
FREETMPS;
LEAVE;
PL_rs = old_rs;
sv_setsv(get_sv("/", GV_ADD), old_rs);
c->in_callback--;
}
static int
psgix_io_svt_get (pTHX_ SV *sv, MAGIC *mg)
{
dSP;
struct feer_conn *c = sv_2feer_conn(mg->mg_obj);
trace("invoking psgix.io magic for fd=%d\n", c->fd);
sv_unmagic(sv, PERL_MAGIC_ext);
ENTER;
SAVETMPS;
PUSHMARK(SP);
XPUSHs(sv);
mXPUSHs(newSViv(c->fd));
PUTBACK;
call_pv("Feersum::Connection::_raw", G_VOID|G_DISCARD|G_EVAL);
SPAGAIN;
if (unlikely(SvTRUE(ERRSV))) {
call_died(aTHX_ c, "psgix.io magic");
}
else {
SV *io_glob = SvRV(sv);
GvSV(io_glob) = newRV_inc(c->self);
// Put whatever remainder data into the socket buffer.
// Optimizes for the websocket case.
//
// TODO: For keepalive support the opposite operation is required;
// pull the data out of the socket buffer and back into feersum.
if (likely(c->rbuf && SvOK(c->rbuf) && SvCUR(c->rbuf))) {
STRLEN rbuf_len;
const char *rbuf_ptr = SvPV(c->rbuf, rbuf_len);
IO *io = GvIOp(io_glob);
assert(io != NULL);
PerlIO_unread(IoIFP(io), (const void *)rbuf_ptr, rbuf_len);
sv_setpvs(c->rbuf, "");
}
stop_read_watcher(c);
stop_read_timer(c);
// don't stop write watcher in case there's outstanding data.
}
PUTBACK;
FREETMPS;
LEAVE;
return 0;
}
MODULE = Feersum PACKAGE = Feersum
PROTOTYPES: ENABLE
void
set_server_name_and_port(SV *self, SV *name, SV *port)
PPCODE:
{
if (feer_server_name)
SvREFCNT_dec(feer_server_name);
feer_server_name = newSVsv(name);
SvREADONLY_on(feer_server_name);
if (feer_server_port)
SvREFCNT_dec(feer_server_port);
feer_server_port = newSVsv(port);
SvREADONLY_on(feer_server_port);
}
void
accept_on_fd(SV *self, int fd)
PPCODE:
{
struct sockaddr_storage addr;
socklen_t addr_len = sizeof(addr);
if (getsockname(fd, (struct sockaddr*)&addr, &addr_len) == -1) perror("getsockname");
switch (addr.ss_family) {
case AF_INET:
case AF_INET6:
is_tcp = 1;
#ifdef TCP_DEFER_ACCEPT
trace("going to defer accept on %d\n",fd);
if (setsockopt(fd, IPPROTO_TCP, TCP_DEFER_ACCEPT, &(int){1}, sizeof(int)) < 0)
perror("setsockopt TCP_DEFER_ACCEPT");
#endif
break;
#ifdef AF_UNIX
case AF_UNIX:
SvREFCNT_dec(request_cb_cv);
request_cb_cv = newSVsv(cb); // copy so 5.8.7 overload magic sticks.
request_cb_is_psgi = ix;
trace("assigned %s request handler %p\n",
request_cb_is_psgi?"PSGI":"Feersum", request_cb_cv);
}
void
graceful_shutdown (SV *self, SV *cb)
PROTOTYPE: $&
PPCODE:
{
if (!IsCodeRef(cb))
croak("must supply a code reference");
if (unlikely(shutting_down))
croak("already shutting down");
shutdown_cb_cv = newSVsv(cb);
trace("shutting down, handler=%p, active=%d\n", SvRV(cb), active_conns);
shutting_down = 1;
ev_io_stop(feersum_ev_loop, &accept_w);
close(accept_w.fd);
if (active_conns <= 0) {
trace("shutdown is immediate\n");
dSP;
ENTER;
SAVETMPS;
PUSHMARK(SP);
call_sv(shutdown_cb_cv, G_EVAL|G_VOID|G_DISCARD|G_NOARGS|G_KEEPERR);
PUTBACK;
trace3("called shutdown handler\n");
SvREFCNT_dec(shutdown_cb_cv);
shutdown_cb_cv = NULL;
FREETMPS;
LEAVE;
}
}
double
read_timeout (SV *self, ...)
PROTOTYPE: $;$
PREINIT:
double new_read_timeout = 0.0;
CODE:
{
if (items > 1) {
new_read_timeout = SvNV(ST(1));
if (!(new_read_timeout > 0.0)) {
croak("must set a positive (non-zero) value for the timeout");
}
trace("set timeout %f\n", new_read_timeout);
read_timeout = new_read_timeout;
}
RETVAL = read_timeout;
}
OUTPUT:
RETVAL
void
set_keepalive (SV *self, SV *set)
PPCODE:
{
trace("set keepalive %d\n", SvTRUE(set));
is_keepalive = SvTRUE(set);
}
unsigned int
max_connection_reqs (SV *self, ...)
PROTOTYPE: $;$
PREINIT:
unsigned int new_max_connection_reqs = 0;
CODE:
{
if (items > 1) {
new_max_connection_reqs = SvIV(ST(1));
if (!(new_max_connection_reqs >= 0)) {
croak("must set a positive value");
}
trace("set max requests per connection %d\n", new_max_connection_reqs);
max_connection_reqs = new_max_connection_reqs;
}
RETVAL = max_connection_reqs;
}
OUTPUT:
RETVAL
void
DESTROY (SV *self)
PPCODE:
{
trace3("DESTROY server\n");
if (request_cb_cv)
SvREFCNT_dec(request_cb_cv);
}
MODULE = Feersum PACKAGE = Feersum::Connection::Handle
PROTOTYPES: ENABLE
int
fileno (feer_conn_handle *hdl)
CODE:
RETVAL = c->fd;
OUTPUT:
RETVAL
void
DESTROY (SV *self)
ALIAS:
Feersum::Connection::Reader::DESTROY = 1
Feersum::Connection::Writer::DESTROY = 2
PPCODE:
{
feer_conn_handle *hdl = sv_2feer_conn_handle(self, 0);
if (hdl == NULL) {
trace3("DESTROY handle (closed) class=%s\n",
HvNAME(SvSTASH(SvRV(self))));
}
else {
struct feer_conn *c = (struct feer_conn *)hdl;
trace3("DESTROY handle fd=%d, class=%s\n", c->fd,
HvNAME(SvSTASH(SvRV(self))));
if (ix == 2) // only close the writer on destruction
PROTOTYPE: $
CODE:
RETVAL = SvREFCNT_inc_simple_NN(feersum_env_addr(aTHX_ c));
OUTPUT:
RETVAL
SV *
remote_port (struct feer_conn *c)
PROTOTYPE: $
CODE:
RETVAL = SvREFCNT_inc_simple_NN(feersum_env_port(aTHX_ c));
OUTPUT:
RETVAL
ssize_t
content_length (struct feer_conn *c)
PROTOTYPE: $
CODE:
RETVAL = feersum_env_content_length(aTHX_ c);
OUTPUT:
RETVAL
SV *
input (struct feer_conn *c)
PROTOTYPE: $
CODE:
if (likely(c->expected_cl > 0)) {
RETVAL = new_feer_conn_handle(aTHX_ c, 0);
} else {
RETVAL = &PL_sv_undef;
}
OUTPUT:
RETVAL
SV *
headers (struct feer_conn *c, int norm = 0)
PROTOTYPE: $;$
CODE:
struct feer_req *r = c->req;
RETVAL = newRV_noinc((SV*)feersum_env_headers(aTHX_ r, norm));
OUTPUT:
RETVAL
SV *
header (struct feer_conn *c, SV *name)
PROTOTYPE: $$
CODE:
struct feer_req *r = c->req;
RETVAL = feersum_env_header(aTHX_ r, name);
OUTPUT:
RETVAL
int
fileno (struct feer_conn *c)
CODE:
RETVAL = c->fd;
OUTPUT:
RETVAL
bool
is_keepalive (struct feer_conn *c)
CODE:
RETVAL = c->is_keepalive;
OUTPUT:
RETVAL
SV*
response_guard (struct feer_conn *c, ...)
PROTOTYPE: $;$
CODE:
RETVAL = feersum_conn_guard(aTHX_ c, (items == 2) ? ST(1) : NULL);
OUTPUT:
RETVAL
void
DESTROY (struct feer_conn *c)
PPCODE:
{
int i;
trace("DESTROY connection fd=%d c=%p\n", c->fd, c);
if (likely(c->rbuf)) SvREFCNT_dec(c->rbuf);
if (c->wbuf_rinq) {
struct iomatrix *m;
while ((m = (struct iomatrix *)rinq_shift(&c->wbuf_rinq)) != NULL) {
for (i=0; i < m->count; i++) {
if (m->sv[i]) SvREFCNT_dec(m->sv[i]);
}
Safefree(m);
}
}
if (likely(c->req)) {
if (c->req->buf) SvREFCNT_dec(c->req->buf);
if (likely(c->req->path)) SvREFCNT_dec(c->req->path);
if (likely(c->req->query)) SvREFCNT_dec(c->req->query);
if (likely(c->req->addr)) SvREFCNT_dec(c->req->addr);
if (likely(c->req->port)) SvREFCNT_dec(c->req->port);
Safefree(c->req);
}
if (likely(c->sa)) free(c->sa);
safe_close_conn(c, "close at destruction");
if (c->poll_write_cb) SvREFCNT_dec(c->poll_write_cb);
if (c->ext_guard) SvREFCNT_dec(c->ext_guard);
active_conns--;
if (unlikely(shutting_down && active_conns <= 0)) {
ev_idle_stop(feersum_ev_loop, &ei);
ev_prepare_stop(feersum_ev_loop, &ep);
ev_check_stop(feersum_ev_loop, &ec);
trace3("... was last conn, going to try shutdown\n");
if (shutdown_cb_cv) {
PUSHMARK(SP);
call_sv(shutdown_cb_cv, G_EVAL|G_VOID|G_DISCARD|G_NOARGS|G_KEEPERR);
PUTBACK;
trace3("... ok, called that handler\n");
( run in 1.646 second using v1.01-cache-2.11-cpan-364913b4093 )