EV-Pg
view release on metacpan or search on metacpan
#include "EXTERN.h"
#include "perl.h"
#include "XSUB.h"
#include "EVAPI.h"
#include <libpq-fe.h>
#include <limits.h>
#include <errno.h>
typedef struct ev_pg_s ev_pg_t;
typedef struct ev_pg_cb_s ev_pg_cb_t;
typedef ev_pg_t* EV__Pg;
typedef struct ev_loop* EV__Loop;
#define EV_PG_MAGIC 0xDEADBEEF
#define EV_PG_FREED 0xFEEDFACE
struct ev_pg_s {
unsigned int magic;
struct ev_loop *loop;
SV *loop_sv; /* ref to the EV::Loop object; keeps the ev_loop alive for our lifetime */
PGconn *conn;
ev_io rio, wio;
int reading, writing;
int rio_unref;
int fd;
int connecting;
char *conninfo;
ev_pg_cb_t *cb_head, *cb_tail;
int pending_count;
int copy_mode;
int draining_single_row;
int draining_copy;
int skip_results;
int meta_fresh;
PGresult *pending_result;
ev_pg_cb_t *delivering_cbt;
PGconn *conn_to_finish;
SV *on_connect;
SV *on_error;
SV *on_notify;
SV *on_notice;
SV *on_drain;
int callback_depth;
int keep_alive;
HV *last_error_fields;
HV *last_result_meta;
PGresult *meta_res;
FILE *trace_fp;
#ifdef LIBPQ_HAS_ASYNC_CANCEL
PGcancelConn *cancel_conn;
ev_io cancel_rio, cancel_wio;
int cancel_reading, cancel_writing;
int cancel_fd;
SV *cancel_cb;
#endif
};
struct ev_pg_cb_s {
SV *cb;
ev_pg_cb_t *next;
int is_pipeline_sync;
int is_describe;
};
static void connect_poll_cb(EV_P_ ev_io *w, int revents);
static void io_read_cb(EV_P_ ev_io *w, int revents);
static void io_write_cb(EV_P_ ev_io *w, int revents);
static void start_reading(ev_pg_t *self);
static void stop_reading(ev_pg_t *self);
static void start_writing(ev_pg_t *self);
static void stop_writing(ev_pg_t *self);
static void drain_notifies(ev_pg_t *self);
static int deliver_result(ev_pg_t *self, PGresult *res);
static void process_results(ev_pg_t *self);
static void check_flush(ev_pg_t *self);
static void emit_error(ev_pg_t *self, const char *msg);
static int cleanup_connection(ev_pg_t *self);
static HV* build_error_fields(PGresult *res);
static HV* build_result_meta(PGresult *res);
static int cancel_pending(ev_pg_t *self, const char *errmsg);
static int check_destroyed(ev_pg_t *self);
static int handle_conn_loss(ev_pg_t *self);
#define REQUIRE_CONN(self) \
if (NULL == (self)->conn || (self)->connecting) croak("not connected")
#define REQUIRE_CB(cb) \
if (!(SvROK(cb) && SvTYPE(SvRV(cb)) == SVt_PVCV)) \
croak("callback must be a CODE reference")
#define CALL_SV_GUARDED(sv, label) \
STMT_START { \
SV *_guarded_sv = SvREFCNT_inc_simple_NN(sv); \
sv_setpvs(ERRSV, ""); \
call_sv(_guarded_sv, G_DISCARD | G_EVAL); \
if (SvTRUE(ERRSV)) \
warn("EV::Pg: exception in " label ": %s", SvPV_nolen(ERRSV)); \
SvREFCNT_dec(_guarded_sv); \
} STMT_END
#define RELEASE_HANDLER(slot) \
STMT_START { if (NULL != (slot)) { SvREFCNT_dec(slot); (slot) = NULL; } } STMT_END
#define STORE_LAST_HV(slot, val) \
STMT_START { \
HV *_new_hv = (val); \
if ((slot)) SvREFCNT_dec((SV*)(slot)); \
(slot) = _new_hv; \
} STMT_END
#define RELEASE_LAST_HV(slot) \
STMT_START { \
if ((slot)) { SvREFCNT_dec((SV*)(slot)); (slot) = NULL; } \
} STMT_END
static ev_pg_cb_t *cbt_freelist = NULL;
static ev_pg_cb_t* alloc_cbt(void) {
ev_pg_cb_t *cbt;
if (cbt_freelist) {
cbt = cbt_freelist;
cbt_freelist = cbt->next;
} else {
Newx(cbt, 1, ev_pg_cb_t);
}
return cbt;
}
static void release_cbt(ev_pg_cb_t *cbt) {
cbt->next = cbt_freelist;
cbt_freelist = cbt;
}
static void start_reading(ev_pg_t *self) {
if (!self->reading && self->fd >= 0) {
ev_io_start(self->loop, &self->rio);
self->reading = 1;
}
}
static void stop_reading(ev_pg_t *self) {
if (self->reading) {
if (self->rio_unref) {
ev_ref(self->loop);
self->rio_unref = 0;
}
ev_io_stop(self->loop, &self->rio);
self->reading = 0;
}
}
/* Call after pending_count changes or connection completes.
* Keeps the read watcher unref'd when idle (no pending queries,
* not connecting) so EV::run can exit. */
static void update_idle_ref(ev_pg_t *self) {
int want_unref;
if (NULL == self->loop) return;
want_unref = self->reading && !self->connecting
&& self->pending_count == 0 && !self->copy_mode
&& !self->draining_single_row && !self->draining_copy
&& !self->keep_alive;
if (want_unref && !self->rio_unref) {
ev_unref(self->loop);
self->rio_unref = 1;
}
else if (!want_unref && self->rio_unref) {
ev_ref(self->loop);
self->rio_unref = 0;
}
}
static void start_writing(ev_pg_t *self) {
if (!self->writing && self->fd >= 0) {
ev_io_start(self->loop, &self->wio);
self->writing = 1;
}
}
static void stop_writing(ev_pg_t *self) {
if (self->writing) {
ev_io_stop(self->loop, &self->wio);
self->writing = 0;
}
}
#ifdef LIBPQ_HAS_ASYNC_CANCEL
static void cancel_poll_cb(EV_P_ ev_io *w, int revents);
static void start_cancel_reading(ev_pg_t *self) {
if (!self->cancel_reading && self->cancel_fd >= 0) {
ev_io_start(self->loop, &self->cancel_rio);
self->cancel_reading = 1;
}
}
static void stop_cancel_reading(ev_pg_t *self) {
if (self->cancel_reading) {
ev_io_stop(self->loop, &self->cancel_rio);
self->cancel_reading = 0;
}
}
static void start_cancel_writing(ev_pg_t *self) {
if (!self->cancel_writing && self->cancel_fd >= 0) {
ev_io_start(self->loop, &self->cancel_wio);
self->cancel_writing = 1;
}
}
static void stop_cancel_writing(ev_pg_t *self) {
if (self->cancel_writing) {
ev_io_stop(self->loop, &self->cancel_wio);
self->cancel_writing = 0;
}
}
static void cleanup_cancel(ev_pg_t *self) {
SV *cb;
stop_cancel_reading(self);
stop_cancel_writing(self);
self->rio.data = (void *)self;
ev_io_init(&self->wio, connect_poll_cb, self->fd, EV_WRITE);
self->wio.data = (void *)self;
start_writing(self);
}
static void begin_connect(ev_pg_t *self, const char *conninfo, const char *what) {
self->conn = PQconnectStart(conninfo);
setup_new_conn(self, what);
}
static const char** marshal_params(AV *params, int nparams,
const char *stack_buf[]) {
const char **pv;
int i;
if (nparams <= 16) {
Zero(stack_buf, nparams, const char *);
pv = stack_buf;
} else {
Newxz(pv, nparams, const char *);
}
for (i = 0; i < nparams; i++) {
SV **svp = av_fetch(params, i, 0);
if (svp) {
SvGETMAGIC(*svp);
if (SvOK(*svp)) {
STRLEN len;
const char *s = SvPV(*svp, len);
if (memchr(s, '\0', len)) {
if (pv != stack_buf) Safefree(pv);
croak("parameter %d contains NUL byte; "
"text-format params cannot contain NULs "
"(use escape_bytea for binary data)", i);
}
pv[i] = s;
}
}
}
return pv;
}
MODULE = EV::Pg PACKAGE = EV::Pg
BOOT:
{
I_EV_API("EV::Pg");
}
EV::Pg
_new(char *class, EV::Loop loop)
CODE:
{
PERL_UNUSED_VAR(class);
Newxz(RETVAL, 1, ev_pg_t);
RETVAL->magic = EV_PG_MAGIC;
RETVAL->loop = loop;
RETVAL->loop_sv = newSVsv(ST(1)); /* keep the loop object alive for our lifetime */
RETVAL->fd = -1;
#ifdef LIBPQ_HAS_ASYNC_CANCEL
RETVAL->cancel_fd = -1;
#endif
}
OUTPUT:
RETVAL
void
DESTROY(EV::Pg self)
CODE:
{
if (self->magic != EV_PG_MAGIC) return;
self->magic = EV_PG_FREED;
/* loop_sv keeps the loop alive for these stops in the normal case;
* during global destruction (PL_dirty) loop teardown order isn't guaranteed. */
stop_reading(self);
stop_writing(self);
if (PL_dirty) {
#ifdef LIBPQ_HAS_ASYNC_CANCEL
stop_cancel_reading(self);
stop_cancel_writing(self);
if (self->cancel_conn) PQcancelFinish(self->cancel_conn);
if (self->cancel_cb) SvREFCNT_dec(self->cancel_cb);
#endif
if (self->pending_result) PQclear(self->pending_result);
if (self->meta_res) PQclear(self->meta_res);
if (self->trace_fp) {
if (self->conn) PQuntrace(self->conn);
fclose(self->trace_fp);
}
if (self->conn) PQfinish(self->conn);
if (self->conn_to_finish) PQfinish(self->conn_to_finish);
Safefree(self->conninfo);
RELEASE_HANDLER(self->on_connect);
RELEASE_HANDLER(self->on_error);
RELEASE_HANDLER(self->on_notify);
RELEASE_HANDLER(self->on_notice);
RELEASE_HANDLER(self->on_drain);
RELEASE_LAST_HV(self->last_error_fields);
RELEASE_LAST_HV(self->last_result_meta);
while (self->cb_head) {
ev_pg_cb_t *cbt = self->cb_head;
self->cb_head = cbt->next;
if (cbt->cb) SvREFCNT_dec(cbt->cb);
Safefree(cbt);
}
self->cb_tail = NULL;
if (self->loop_sv) SvREFCNT_dec(self->loop_sv);
Safefree(self);
return;
}
if (self->pending_result) {
PQclear(self->pending_result);
self->pending_result = NULL;
}
if (self->meta_res) {
PQclear(self->meta_res);
self->meta_res = NULL;
}
CLEANUP_CANCEL(self);
{
PGconn *conn = self->conn;
self->conn = NULL;
self->loop = NULL;
self->fd = -1;
if (self->trace_fp && conn) PQuntrace(conn);
if (conn) PQfinish(conn);
if (self->conn_to_finish) {
PQfinish(self->conn_to_finish);
self->conn_to_finish = NULL;
}
SV*
ssl_attribute(EV::Pg self, const char *name)
CODE:
{
RETVAL = conn_str_or_undef(self->conn ? PQsslAttribute(self->conn, name) : NULL);
}
OUTPUT:
RETVAL
SV*
escape_literal(EV::Pg self, SV *str)
PREINIT:
STRLEN len;
const char *s;
char *escaped;
CODE:
{
REQUIRE_CONN(self);
s = SvPV(str, len);
escaped = PQescapeLiteral(self->conn, s, len);
if (NULL == escaped) {
croak("PQescapeLiteral failed: %s", PQerrorMessage(self->conn));
}
RETVAL = newSVpv(escaped, 0);
PQfreemem(escaped);
}
OUTPUT:
RETVAL
SV*
escape_identifier(EV::Pg self, SV *str)
PREINIT:
STRLEN len;
const char *s;
char *escaped;
CODE:
{
REQUIRE_CONN(self);
s = SvPV(str, len);
escaped = PQescapeIdentifier(self->conn, s, len);
if (NULL == escaped) {
croak("PQescapeIdentifier failed: %s", PQerrorMessage(self->conn));
}
RETVAL = newSVpv(escaped, 0);
PQfreemem(escaped);
}
OUTPUT:
RETVAL
int
pending_count(EV::Pg self)
CODE:
{
RETVAL = self->pending_count;
}
OUTPUT:
RETVAL
int
keep_alive(EV::Pg self, ...)
CODE:
{
if (items > 1) {
self->keep_alive = SvTRUE(ST(1)) ? 1 : 0;
update_idle_ref(self);
}
RETVAL = self->keep_alive;
}
OUTPUT:
RETVAL
void
skip_pending(EV::Pg self)
CODE:
{
PGconn *entry_conn = self->conn;
int skipped = cancel_pending(self, "skipped");
if (self->magic != EV_PG_MAGIC) {
check_destroyed(self);
return;
}
/* Only credit if the same connection is still in place: a "skipped"
* callback may have torn it down (finish/reset/connect), taking its
* in-flight sequences â and skip_results â with it. */
if (skipped > 0 && self->conn == entry_conn)
self->skip_results += skipped;
check_destroyed(self);
}
int
lib_version(char *class)
CODE:
{
PERL_UNUSED_VAR(class);
RETVAL = PQlibVersion();
}
OUTPUT:
RETVAL
SV*
conninfo_parse(char *class, const char *conninfo)
CODE:
{
PQconninfoOption *opts;
char *errmsg = NULL;
PERL_UNUSED_VAR(class);
opts = PQconninfoParse(conninfo, &errmsg);
if (!opts) {
if (errmsg) {
SV *errsv = newSVpv(errmsg, 0);
SAVEFREESV(errsv);
PQfreemem(errmsg);
croak("PQconninfoParse failed: %s", SvPV_nolen(errsv));
}
croak("PQconninfoParse failed");
}
RETVAL = conninfo_opts_to_hv(opts);
}
OUTPUT:
RETVAL
SV*
cancel(EV::Pg self)
PREINIT:
PGcancel *cn;
( run in 1.333 second using v1.01-cache-2.11-cpan-941387dca55 )