Feersum
view release on metacpan or search on metacpan
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;
INLINE_UNLESS_DEBUG static void
feersum_set_path_and_query(pTHX_ struct feer_req *r)
{
const char *qpos = r->uri;
while (*qpos != '?' && qpos < r->uri + r->uri_len) qpos++;
if (*qpos == '?') {
r->path = newSVpvn(r->uri, (qpos - r->uri));
qpos++;
r->query = newSVpvn(qpos, r->uri_len - (qpos - r->uri));
} else {
r->path = feersum_env_uri(aTHX_ r);
r->query = newSVpvs("");
}
uri_decode_sv(r->path);
}
INLINE_UNLESS_DEBUG static SV*
feersum_env_path(pTHX_ struct feer_req *r)
{
if (unlikely(!r->path)) feersum_set_path_and_query(aTHX_ r);
return r->path;
}
INLINE_UNLESS_DEBUG static SV*
feersum_env_query(pTHX_ struct feer_req *r)
{
if (unlikely(!r->query)) feersum_set_path_and_query(aTHX_ r);
return r->query;
}
INLINE_UNLESS_DEBUG static SV*
feersum_env_addr(pTHX_ struct feer_conn *c)
{
struct feer_req *r = c->req;
if (unlikely(!r->addr)) feersum_set_remote_info(aTHX_ r, c->sa);
return r->addr;
}
INLINE_UNLESS_DEBUG static SV*
feersum_env_port(pTHX_ struct feer_conn *c)
{
struct feer_req *r = c->req;
if (unlikely(!r->port)) feersum_set_remote_info(aTHX_ r, c->sa);
return r->port;
}
static void
feersum_init_tmpl_env(pTHX)
{
HV *e;
e = newHV();
// constants
hv_stores(e, "psgi.version", newRV((SV*)psgi_ver));
hv_stores(e, "psgi.url_scheme", newSVpvs("http"));
hv_stores(e, "psgi.run_once", &PL_sv_no);
hv_stores(e, "psgi.nonblocking", &PL_sv_yes);
hv_stores(e, "psgi.multithread", &PL_sv_no);
hv_stores(e, "psgi.multiprocess", &PL_sv_no);
hv_stores(e, "psgi.streaming", &PL_sv_yes);
hv_stores(e, "psgi.errors", newRV((SV*)PL_stderrgv));
hv_stores(e, "psgix.input.buffered", &PL_sv_yes);
hv_stores(e, "psgix.output.buffered", &PL_sv_yes);
hv_stores(e, "psgix.body.scalar_refs", &PL_sv_yes);
hv_stores(e, "psgix.output.guard", &PL_sv_yes);
hv_stores(e, "SCRIPT_NAME", newSVpvs(""));
// placeholders that get defined for every request
hv_stores(e, "SERVER_PROTOCOL", &PL_sv_undef);
hv_stores(e, "SERVER_NAME", &PL_sv_undef);
hv_stores(e, "SERVER_PORT", &PL_sv_undef);
hv_stores(e, "REQUEST_URI", &PL_sv_undef);
hv_stores(e, "REQUEST_METHOD", &PL_sv_undef);
hv_stores(e, "PATH_INFO", &PL_sv_undef);
hv_stores(e, "REMOTE_ADDR", &PL_sv_placeholder);
hv_stores(e, "REMOTE_PORT", &PL_sv_placeholder);
// defaults that get changed for some requests
hv_stores(e, "psgi.input", &PL_sv_undef);
hv_stores(e, "CONTENT_LENGTH", newSViv(0));
hv_stores(e, "QUERY_STRING", newSVpvs(""));
// anticipated headers
hv_stores(e, "CONTENT_TYPE", &PL_sv_placeholder);
hv_stores(e, "HTTP_HOST", &PL_sv_placeholder);
hv_stores(e, "HTTP_USER_AGENT", &PL_sv_placeholder);
hv_stores(e, "HTTP_ACCEPT", &PL_sv_placeholder);
hv_stores(e, "HTTP_ACCEPT_LANGUAGE", &PL_sv_placeholder);
hv_stores(e, "HTTP_ACCEPT_CHARSET", &PL_sv_placeholder);
hv_stores(e, "HTTP_KEEP_ALIVE", &PL_sv_placeholder);
hv_stores(e, "HTTP_CONNECTION", &PL_sv_placeholder);
hv_stores(e, "HTTP_REFERER", &PL_sv_placeholder);
hv_stores(e, "HTTP_COOKIE", &PL_sv_placeholder);
hv_stores(e, "HTTP_IF_MODIFIED_SINCE", &PL_sv_placeholder);
hv_stores(e, "HTTP_IF_NONE_MATCH", &PL_sv_placeholder);
hv_stores(e, "HTTP_CACHE_CONTROL", &PL_sv_placeholder);
hv_stores(e, "psgix.io", &PL_sv_placeholder);
feersum_tmpl_env = e;
}
static HV*
feersum_env(pTHX_ struct feer_conn *c)
{
HV *e;
SV **hsv;
int i,j;
struct feer_req *r = c->req;
if (unlikely(!feersum_tmpl_env))
feersum_init_tmpl_env(aTHX);
e = newHVhv(feersum_tmpl_env);
trace("generating header (fd %d) %.*s\n",
c->fd, (int)r->uri_len, r->uri);
hv_stores(e, "SERVER_NAME", SvREFCNT_inc_simple(feer_server_name));
hv_stores(e, "SERVER_PORT", newSVsv_nomg(feer_server_port));
hv_stores(e, "REQUEST_URI", feersum_env_uri(aTHX_ r));
char *k = kbuf;\
for (j = 0; j < hdr->name_len; j++) { char n = hdr->name[j]; *k++ = _str; }\
if (unlikely(kbuflen < hdr->name_len)) { kbuflen = hdr->name_len; kbuf = Renew(kbuf, kbuflen, char); }\
SV** val = hv_fetch(e, kbuf, hdr->name_len, 1);\
if (unlikely(SvPOK(*val))) {\
sv_catpvn(*val, ", ", 2);\
sv_catpvn(*val, hdr->value, hdr->value_len);\
} else {\
sv_setpvn(*val, hdr->value, hdr->value_len);\
}\
}\
break;
INLINE_UNLESS_DEBUG static HV*
feersum_env_headers(pTHX_ struct feer_req *r, int norm)
{
int i; int j; char* n; HV* e;
e = newHV();
SV** val;
char *kbuf;
size_t kbuflen = 64;
Newx(kbuf, kbuflen, char);
switch (norm) {
case HEADER_NORM_SKIP:
COPY_NORM_HEADER(n)
case HEADER_NORM_LOCASE:
COPY_NORM_HEADER(tolower(n))
case HEADER_NORM_UPCASE:
COPY_NORM_HEADER(toupper(n))
case HEADER_NORM_LOCASE_DASH:
COPY_NORM_HEADER((n == '-') ? '_' : tolower(n))
case HEADER_NORM_UPCASE_DASH:
COPY_NORM_HEADER((n == '-') ? '_' : toupper(n))
}
Safefree(kbuf);
return e;
}
INLINE_UNLESS_DEBUG static SV*
feersum_env_header(pTHX_ struct feer_req *r, SV *name)
{
int i;
for (i = 0; i < r->num_headers; i++) {
struct phr_header *hdr = &(r->headers[i]);
if (hdr->name == NULL) continue;
if (unlikely(str_case_eq(SvPVX(name), SvCUR(name), hdr->name, hdr->name_len))) {
return newSVpvn(hdr->value, hdr->value_len);
}
}
return &PL_sv_undef;
}
INLINE_UNLESS_DEBUG static ssize_t
feersum_env_content_length(pTHX_ struct feer_conn *c)
{
return c->expected_cl;
}
static void
feersum_start_response (pTHX_ struct feer_conn *c, SV *message, AV *headers,
int streaming)
{
const char *ptr;
I32 i;
trace("start_response fd=%d streaming=%d\n", c->fd, streaming);
if (unlikely(c->responding != RESPOND_NOT_STARTED))
croak("already responding?!");
change_responding_state(c, streaming ? RESPOND_STREAMING : RESPOND_NORMAL);
if (unlikely(!SvOK(message) || !(SvIOK(message) || SvPOK(message)))) {
croak("Must define an HTTP status code or message");
}
I32 avl = av_len(headers);
if (unlikely(avl+1 % 2 == 1)) {
croak("expected even-length array, got %d", avl+1);
}
// int or 3 chars? use a stock message
UV code = 0;
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)) {
sv_setpvf(cl_sv, "Content-Length: %"Sz_uf"" CRLFx2, (Sz)RETVAL);
update_wbuf_placeholder(c, cl_sv, cl_iov);
}
change_responding_state(c, RESPOND_SHUTDOWN);
conn_write_ready(c);
return RETVAL;
}
static void
feersum_start_psgi_streaming(pTHX_ struct feer_conn *c, SV *streamer)
{
dSP;
ENTER;
SAVETMPS;
PUSHMARK(SP);
mXPUSHs(feer_conn_2sv(c));
XPUSHs(streamer);
PUTBACK;
call_method("_initiate_streaming_psgi", G_DISCARD|G_EVAL|G_VOID);
SPAGAIN;
if (unlikely(SvTRUE(ERRSV))) {
call_died(aTHX_ c, "PSGI stream initiator");
}
PUTBACK;
FREETMPS;
LEAVE;
}
static void
feersum_handle_psgi_response(
pTHX_ struct feer_conn *c, SV *ret, bool can_recurse)
{
if (unlikely(!SvOK(ret) || !SvROK(ret))) {
sv_setpvs(ERRSV, "Invalid PSGI response (expected reference)");
call_died(aTHX_ c, "PSGI request");
return;
}
if (SvOK(ret) && unlikely(!IsArrayRef(ret))) {
if (likely(can_recurse)) {
trace("PSGI response non-array, c=%p ret=%p\n", c, ret);
feersum_start_psgi_streaming(aTHX_ c, ret);
}
else {
sv_setpvs(ERRSV, "PSGI attempt to recurse in a streaming callback");
call_died(aTHX_ c, "PSGI request");
}
return;
}
AV *psgi_triplet = (AV*)SvRV(ret);
if (unlikely(av_len(psgi_triplet)+1 != 3)) {
sv_setpvs(ERRSV, "Invalid PSGI array response (expected triplet)");
call_died(aTHX_ c, "PSGI request");
return;
}
trace("PSGI response triplet, c=%p av=%p\n", c, psgi_triplet);
// we know there's three elems so *should* be safe to de-ref
SV *msg = *(av_fetch(psgi_triplet,0,0));
SV *hdrs = *(av_fetch(psgi_triplet,1,0));
SV *body = *(av_fetch(psgi_triplet,2,0));
AV *headers;
if (IsArrayRef(hdrs))
headers = (AV*)SvRV(hdrs);
else {
sv_setpvs(ERRSV, "PSGI Headers must be an array-ref");
call_died(aTHX_ c, "PSGI request");
return;
}
if (likely(IsArrayRef(body))) {
feersum_start_response(aTHX_ c, msg, headers, 0);
feersum_write_whole_body(aTHX_ c, body);
}
else if (likely(SvROK(body))) { // probaby an IO::Handle-like object
feersum_start_response(aTHX_ c, msg, headers, 1);
c->poll_write_cb = newSVsv(body);
c->poll_write_cb_is_io_handle = 1;
conn_write_ready(c);
}
else {
sv_setpvs(ERRSV, "Expected PSGI array-ref or IO::Handle-like body");
call_died(aTHX_ c, "PSGI request");
return;
}
}
static int
feersum_close_handle (pTHX_ struct feer_conn *c, bool is_writer)
{
int RETVAL;
if (is_writer) {
trace("close writer fd=%d, c=%p, refcnt=%d\n", c->fd, c, SvREFCNT(c->self));
if (c->poll_write_cb) {
SvREFCNT_dec(c->poll_write_cb);
c->poll_write_cb = NULL;
}
if (c->responding < RESPOND_SHUTDOWN) {
finish_wbuf(c);
conn_write_ready(c);
change_responding_state(c, RESPOND_SHUTDOWN);
}
RETVAL = 1;
buf_ptr = SvPV(buf, buf_len);
if (likely(c->rbuf))
src_ptr = SvPV(c->rbuf, src_len);
if (unlikely(len < 0))
len = src_len;
if (unlikely(offset < 0))
offset = (-offset >= c->received_cl) ? 0 : c->received_cl + offset;
if (unlikely(len + offset > src_len))
len = src_len - offset;
trace("read fd=%d : normalized len=%"Sz_uf" off=%"Ssz_df" src_len=%"Sz_uf"\n",
c->fd, (Sz)len, (Ssz)offset, (Sz)src_len);
if (unlikely(!c->rbuf || src_len == 0 || offset >= c->received_cl)) {
trace2("rbuf empty during read %d\n", c->fd);
if (c->receiving == RECEIVE_SHUTDOWN) {
XSRETURN_IV(0);
}
else {
errno = EAGAIN;
XSRETURN_UNDEF;
}
}
if (likely(len == src_len && offset == 0)) {
trace2("appending entire rbuf fd=%d\n", c->fd);
sv_2mortal(c->rbuf); // allow pv to be stolen
if (likely(buf_len == 0)) {
sv_setsv(buf, c->rbuf);
}
else {
sv_catsv(buf, c->rbuf);
}
c->rbuf = NULL;
}
else {
src_ptr += offset;
trace2("appending partial rbuf fd=%d len=%"Sz_uf" off=%"Ssz_df" ptr=%p\n",
c->fd, len, offset, src_ptr);
SvGROW(buf, SvCUR(buf) + len);
sv_catpvn(buf, src_ptr, len);
if (likely(items == 3)) {
// there wasn't an offset param, throw away beginning
sv_chop(c->rbuf, SvPVX(c->rbuf) + len);
}
}
XSRETURN_IV(len);
}
STRLEN
write (feer_conn_handle *hdl, ...)
PROTOTYPE: $;$
CODE:
{
if (unlikely(c->responding != RESPOND_STREAMING))
croak("can only call write in streaming mode");
SV *body = (items == 2) ? ST(1) : &PL_sv_undef;
if (unlikely(!body || !SvOK(body)))
XSRETURN_IV(0);
trace("write fd=%d c=%p, body=%p\n", c->fd, c, body);
if (SvROK(body)) {
SV *refd = SvRV(body);
if (SvOK(refd) && SvPOK(refd)) {
body = refd;
}
else {
croak("body must be a scalar, scalar ref or undef");
}
}
(void)SvPV(body, RETVAL);
if (c->is_http11)
add_chunk_sv_to_wbuf(c, body);
else
add_sv_to_wbuf(c, body);
conn_write_ready(c);
}
OUTPUT:
RETVAL
void
write_array (feer_conn_handle *hdl, AV *abody)
PROTOTYPE: $$
PPCODE:
{
if (unlikely(c->responding != RESPOND_STREAMING))
croak("can only call write in streaming mode");
trace("write_array fd=%d c=%p, abody=%p\n", c->fd, c, abody);
I32 amax = av_len(abody);
int i;
if (c->is_http11) {
for (i=0; i<=amax; i++) {
SV *sv = fetch_av_normal(aTHX_ abody, i);
if (likely(sv)) add_chunk_sv_to_wbuf(c, sv);
}
}
else {
for (i=0; i<=amax; i++) {
SV *sv = fetch_av_normal(aTHX_ abody, i);
if (likely(sv)) add_sv_to_wbuf(c, sv);
}
}
conn_write_ready(c);
}
int
seek (feer_conn_handle *hdl, ssize_t offset, ...)
PROTOTYPE: $$;$
CODE:
{
int whence = SEEK_CUR;
if (items == 3 && SvOK(ST(2)) && SvIOK(ST(2)))
whence = SvIV(ST(2));
trace("seek fd=%d offset=%"Ssz_df" whence=%d\n", c->fd, offset, whence);
if (unlikely(!c->rbuf)) {
// handle is effectively "closed"
RETVAL = 0;
}
else if (offset == 0) {
RETVAL = 1; // stay put for any whence
}
else if (offset > 0 && (whence == SEEK_CUR || whence == SEEK_SET)) {
STRLEN len;
const char *str = SvPV_const(c->rbuf, len);
if (offset > len)
offset = len;
sv_chop(c->rbuf, str + offset);
RETVAL = 1;
}
else if (offset < 0 && whence == SEEK_END) {
STRLEN len;
const char *str = SvPV_const(c->rbuf, len);
offset += len; // can't be > len since block is offset<0
if (offset == 0) {
RETVAL = 1; // no-op, but OK
}
else if (offset > 0) {
sv_chop(c->rbuf, str + offset);
RETVAL = 1;
}
else {
// past beginning of string
OUTPUT:
RETVAL
int
close (feer_conn_handle *hdl)
PROTOTYPE: $
ALIAS:
Feersum::Connection::Reader::close = 1
Feersum::Connection::Writer::close = 2
CODE:
{
assert(ix);
RETVAL = feersum_close_handle(aTHX_ c, (ix == 2));
SvUVX(hdl_sv) = 0;
}
OUTPUT:
RETVAL
void
_poll_cb (feer_conn_handle *hdl, SV *cb)
PROTOTYPE: $$
ALIAS:
Feersum::Connection::Reader::poll_cb = 1
Feersum::Connection::Writer::poll_cb = 2
PPCODE:
{
if (unlikely(ix < 1 || ix > 2))
croak("can't call _poll_cb directly");
else if (unlikely(ix == 1))
croak("poll_cb for reading not yet supported"); // TODO poll_read_cb
if (c->poll_write_cb != NULL) {
SvREFCNT_dec(c->poll_write_cb);
c->poll_write_cb = NULL;
}
if (!SvOK(cb)) {
trace("unset poll_cb ix=%d\n", ix);
return;
}
else if (unlikely(!IsCodeRef(cb)))
croak("must supply a code reference to poll_cb");
c->poll_write_cb = newSVsv(cb);
conn_write_ready(c);
}
SV*
response_guard (feer_conn_handle *hdl, ...)
PROTOTYPE: $;$
CODE:
RETVAL = feersum_conn_guard(aTHX_ c, (items==2) ? ST(1) : NULL);
OUTPUT:
RETVAL
MODULE = Feersum PACKAGE = Feersum::Connection
PROTOTYPES: ENABLE
SV *
start_streaming (struct feer_conn *c, SV *message, AV *headers)
PROTOTYPE: $$\@
CODE:
feersum_start_response(aTHX_ c, message, headers, 1);
RETVAL = new_feer_conn_handle(aTHX_ c, 1); // RETVAL gets mortalized
OUTPUT:
RETVAL
int
is_http11 (struct feer_conn *c)
CODE:
RETVAL = c->is_http11;
OUTPUT:
RETVAL
size_t
send_response (struct feer_conn *c, SV* message, AV *headers, SV *body)
PROTOTYPE: $$\@$
CODE:
feersum_start_response(aTHX_ c, message, headers, 0);
if (unlikely(!SvOK(body)))
croak("can't send_response with an undef body");
RETVAL = feersum_write_whole_body(aTHX_ c, body);
OUTPUT:
RETVAL
SV*
_continue_streaming_psgi (struct feer_conn *c, SV *psgi_response)
PROTOTYPE: $\@
CODE:
{
AV *av;
int len = 0;
if (IsArrayRef(psgi_response)) {
av = (AV*)SvRV(psgi_response);
len = av_len(av) + 1;
}
if (len == 3) {
// 0 is "don't recurse" (i.e. don't allow another code-ref)
feersum_handle_psgi_response(aTHX_ c, psgi_response, 0);
RETVAL = &PL_sv_undef;
}
else if (len == 2) {
SV *message = *(av_fetch(av,0,0));
SV *headers = *(av_fetch(av,1,0));
if (unlikely(!IsArrayRef(headers)))
croak("PSGI headers must be an array ref");
feersum_start_response(aTHX_ c, message, (AV*)SvRV(headers), 1);
RETVAL = new_feer_conn_handle(aTHX_ c, 1); // RETVAL gets mortalized
}
else {
croak("PSGI response starter expects a 2 or 3 element array-ref");
}
}
OUTPUT:
RETVAL
void
force_http10 (struct feer_conn *c)
PROTOTYPE: $
ALIAS:
force_http11 = 1
PPCODE:
c->is_http11 = ix;
SV *
env (struct feer_conn *c)
PROTOTYPE: $
CODE:
RETVAL = newRV_noinc((SV*)feersum_env(aTHX_ c));
OUTPUT:
RETVAL
SV *
method (struct feer_conn *c)
PROTOTYPE: $
CODE:
struct feer_req *r = c->req;
RETVAL = feersum_env_method(aTHX_ r);
OUTPUT:
RETVAL
SV *
uri (struct feer_conn *c)
PROTOTYPE: $
CODE:
( run in 1.680 second using v1.01-cache-2.11-cpan-5c0b1e786e0 )