EV-Pg

 view release on metacpan or  search on metacpan

Pg.xs  view on Meta::CPAN

#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);

Pg.xs  view on Meta::CPAN

    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;
        }

Pg.xs  view on Meta::CPAN


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 )