2009-02-02 19:22:13 +00:00
|
|
|
/*
|
2009-04-19 20:54:12 +00:00
|
|
|
* Copyright (c) 2007-2009 Niels Provos and Nick Mathewson
|
|
|
|
* Copyright (c) 2002-2006 Niels Provos <provos@citi.umich.edu>
|
2009-02-02 19:22:13 +00:00
|
|
|
* All rights reserved.
|
|
|
|
*
|
|
|
|
* Redistribution and use in source and binary forms, with or without
|
|
|
|
* modification, are permitted provided that the following conditions
|
|
|
|
* are met:
|
|
|
|
* 1. Redistributions of source code must retain the above copyright
|
|
|
|
* notice, this list of conditions and the following disclaimer.
|
|
|
|
* 2. Redistributions in binary form must reproduce the above copyright
|
|
|
|
* notice, this list of conditions and the following disclaimer in the
|
|
|
|
* documentation and/or other materials provided with the distribution.
|
|
|
|
* 3. The name of the author may not be used to endorse or promote products
|
|
|
|
* derived from this software without specific prior written permission.
|
|
|
|
*
|
|
|
|
* THIS SOFTWARE IS PROVIDED BY THE AUTHOR ``AS IS'' AND ANY EXPRESS OR
|
|
|
|
* IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES
|
|
|
|
* OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE DISCLAIMED.
|
|
|
|
* IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR ANY DIRECT, INDIRECT,
|
|
|
|
* INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT
|
|
|
|
* NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
|
|
|
|
* DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
|
|
|
|
* THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
|
|
|
|
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF
|
|
|
|
* THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
|
|
|
*/
|
|
|
|
|
|
|
|
#include <sys/types.h>
|
|
|
|
|
|
|
|
#ifdef HAVE_CONFIG_H
|
2009-02-02 21:24:04 +00:00
|
|
|
#include "event-config.h"
|
2009-02-02 19:22:13 +00:00
|
|
|
#endif
|
|
|
|
|
2009-02-02 21:24:04 +00:00
|
|
|
#ifdef _EVENT_HAVE_SYS_TIME_H
|
2009-02-02 19:22:13 +00:00
|
|
|
#include <sys/time.h>
|
|
|
|
#endif
|
|
|
|
|
|
|
|
#include <errno.h>
|
|
|
|
#include <stdio.h>
|
|
|
|
#include <stdlib.h>
|
|
|
|
#include <string.h>
|
|
|
|
#include <assert.h>
|
2009-02-02 21:24:04 +00:00
|
|
|
#ifdef _EVENT_HAVE_STDARG_H
|
2009-02-02 19:22:13 +00:00
|
|
|
#include <stdarg.h>
|
|
|
|
#endif
|
|
|
|
|
|
|
|
#ifdef WIN32
|
|
|
|
#include <winsock2.h>
|
|
|
|
#endif
|
|
|
|
|
|
|
|
#include "event2/util.h"
|
|
|
|
#include "event2/bufferevent.h"
|
|
|
|
#include "event2/buffer.h"
|
|
|
|
#include "event2/bufferevent_struct.h"
|
|
|
|
#include "event2/event.h"
|
|
|
|
#include "log-internal.h"
|
|
|
|
#include "mm-internal.h"
|
|
|
|
#include "bufferevent-internal.h"
|
|
|
|
#include "util-internal.h"
|
|
|
|
|
|
|
|
/* prototypes */
|
|
|
|
static int be_filter_enable(struct bufferevent *, short);
|
|
|
|
static int be_filter_disable(struct bufferevent *, short);
|
|
|
|
static void be_filter_destruct(struct bufferevent *);
|
|
|
|
|
|
|
|
static void be_filter_readcb(struct bufferevent *, void *);
|
|
|
|
static void be_filter_writecb(struct bufferevent *, void *);
|
2009-05-25 23:11:20 +00:00
|
|
|
static void be_filter_eventcb(struct bufferevent *, short, void *);
|
2009-02-02 19:22:13 +00:00
|
|
|
static int be_filter_flush(struct bufferevent *bufev,
|
|
|
|
short iotype, enum bufferevent_flush_mode mode);
|
2009-05-13 20:37:21 +00:00
|
|
|
static int be_filter_ctrl(struct bufferevent *, enum bufferevent_ctrl_op, union bufferevent_ctrl_data *);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
static void bufferevent_filtered_outbuf_cb(struct evbuffer *buf,
|
2009-04-03 14:27:03 +00:00
|
|
|
const struct evbuffer_cb_info *info, void *arg);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
struct bufferevent_filtered {
|
2009-04-13 03:08:11 +00:00
|
|
|
struct bufferevent_private bev;
|
2009-02-02 19:22:13 +00:00
|
|
|
|
2009-10-16 13:19:57 +00:00
|
|
|
/** The bufferevent that we read/write filtered data from/to. */
|
2009-02-02 19:22:13 +00:00
|
|
|
struct bufferevent *underlying;
|
|
|
|
/** A callback on our outbuf to notice when somebody adds data */
|
|
|
|
struct evbuffer_cb_entry *outbuf_cb;
|
|
|
|
/** True iff we have received an EOF callback from the underlying
|
|
|
|
* bufferevent. */
|
|
|
|
unsigned got_eof;
|
|
|
|
|
|
|
|
/** Function to free context when we're done. */
|
|
|
|
void (*free_context)(void *);
|
|
|
|
/** Input filter */
|
|
|
|
bufferevent_filter_cb process_in;
|
|
|
|
/** Output filter */
|
|
|
|
bufferevent_filter_cb process_out;
|
|
|
|
|
|
|
|
/** User-supplied argument to the filters. */
|
|
|
|
void *context;
|
|
|
|
};
|
|
|
|
|
2009-04-10 15:01:31 +00:00
|
|
|
const struct bufferevent_ops bufferevent_ops_filter = {
|
2009-02-02 19:22:13 +00:00
|
|
|
"filter",
|
|
|
|
evutil_offsetof(struct bufferevent_filtered, bev),
|
|
|
|
be_filter_enable,
|
|
|
|
be_filter_disable,
|
|
|
|
be_filter_destruct,
|
2009-05-25 23:10:23 +00:00
|
|
|
_bufferevent_generic_adj_timeouts,
|
2009-02-02 19:22:13 +00:00
|
|
|
be_filter_flush,
|
2009-05-13 20:37:21 +00:00
|
|
|
be_filter_ctrl,
|
2009-02-02 19:22:13 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
/* Given a bufferevent that's really the bev filter of a bufferevent_filtered,
|
|
|
|
* return that bufferevent_filtered. Returns NULL otherwise.*/
|
|
|
|
static inline struct bufferevent_filtered *
|
|
|
|
upcast(struct bufferevent *bev)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bev_f;
|
|
|
|
if (bev->be_ops != &bufferevent_ops_filter)
|
|
|
|
return NULL;
|
|
|
|
bev_f = (void*)( ((char*)bev) -
|
2009-04-13 03:08:11 +00:00
|
|
|
evutil_offsetof(struct bufferevent_filtered, bev.bev));
|
|
|
|
assert(bev_f->bev.bev.be_ops == &bufferevent_ops_filter);
|
2009-02-02 19:22:13 +00:00
|
|
|
return bev_f;
|
|
|
|
}
|
|
|
|
|
2009-04-13 03:08:11 +00:00
|
|
|
#define downcast(bev_f) (&(bev_f)->bev.bev)
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
/** Return 1 iff bevf's underlying bufferevent's output buffer is at or
|
|
|
|
* over its high watermark such that we should not write to it in a given
|
|
|
|
* flush mode. */
|
|
|
|
static int
|
|
|
|
be_underlying_writebuf_full(struct bufferevent_filtered *bevf,
|
|
|
|
enum bufferevent_flush_mode state)
|
|
|
|
{
|
|
|
|
struct bufferevent *u = bevf->underlying;
|
|
|
|
return state == BEV_NORMAL &&
|
|
|
|
u->wm_write.high &&
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(u->output) >= u->wm_write.high;
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/** Return 1 if our input buffer is at or over its high watermark such that we
|
|
|
|
* should not write to it in a given flush mode. */
|
|
|
|
static int
|
|
|
|
be_readbuf_full(struct bufferevent_filtered *bevf,
|
|
|
|
enum bufferevent_flush_mode state)
|
|
|
|
{
|
|
|
|
struct bufferevent *bufev = downcast(bevf);
|
|
|
|
return state == BEV_NORMAL &&
|
|
|
|
bufev->wm_read.high &&
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bufev->input) >= bufev->wm_read.high;
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/* Filter to use when we're created with a NULL filter. */
|
|
|
|
static enum bufferevent_filter_result
|
2009-05-22 19:11:48 +00:00
|
|
|
be_null_filter(struct evbuffer *src, struct evbuffer *dst, ev_ssize_t lim,
|
2009-02-02 19:22:13 +00:00
|
|
|
enum bufferevent_flush_mode state, void *ctx)
|
|
|
|
{
|
|
|
|
(void)state;
|
2009-02-11 05:17:27 +00:00
|
|
|
if (evbuffer_remove_buffer(src, dst, lim) == 0)
|
2009-02-02 19:22:13 +00:00
|
|
|
return BEV_OK;
|
|
|
|
else
|
|
|
|
return BEV_ERROR;
|
|
|
|
}
|
|
|
|
|
|
|
|
struct bufferevent *
|
|
|
|
bufferevent_filter_new(struct bufferevent *underlying,
|
|
|
|
bufferevent_filter_cb input_filter,
|
|
|
|
bufferevent_filter_cb output_filter,
|
|
|
|
enum bufferevent_options options,
|
|
|
|
void (*free_context)(void *),
|
|
|
|
void *ctx)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bufev_f;
|
2009-04-13 03:17:19 +00:00
|
|
|
enum bufferevent_options tmp_options = options & ~BEV_OPT_THREADSAFE;
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
if (!input_filter)
|
|
|
|
input_filter = be_null_filter;
|
|
|
|
if (!output_filter)
|
|
|
|
output_filter = be_null_filter;
|
|
|
|
|
|
|
|
bufev_f = mm_calloc(1, sizeof(struct bufferevent_filtered));
|
|
|
|
if (!bufev_f)
|
|
|
|
return NULL;
|
|
|
|
|
|
|
|
if (bufferevent_init_common(&bufev_f->bev, underlying->ev_base,
|
2009-04-13 03:17:19 +00:00
|
|
|
&bufferevent_ops_filter, tmp_options) < 0) {
|
2009-02-02 19:22:13 +00:00
|
|
|
mm_free(bufev_f);
|
|
|
|
return NULL;
|
|
|
|
}
|
2009-04-13 03:17:19 +00:00
|
|
|
if (options & BEV_OPT_THREADSAFE) {
|
2009-07-28 04:03:57 +00:00
|
|
|
bufferevent_enable_locking(downcast(bufev_f), NULL);
|
2009-04-13 03:17:19 +00:00
|
|
|
}
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
bufev_f->underlying = underlying;
|
|
|
|
bufev_f->process_in = input_filter;
|
|
|
|
bufev_f->process_out = output_filter;
|
|
|
|
bufev_f->free_context = free_context;
|
|
|
|
bufev_f->context = ctx;
|
|
|
|
|
|
|
|
bufferevent_setcb(bufev_f->underlying,
|
2009-05-25 23:11:20 +00:00
|
|
|
be_filter_readcb, be_filter_writecb, be_filter_eventcb, bufev_f);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
2009-04-13 03:08:11 +00:00
|
|
|
bufev_f->outbuf_cb = evbuffer_add_cb(downcast(bufev_f)->output,
|
2009-02-02 19:22:13 +00:00
|
|
|
bufferevent_filtered_outbuf_cb, bufev_f);
|
|
|
|
|
2009-05-25 23:10:23 +00:00
|
|
|
_bufferevent_init_generic_timeout_cbs(downcast(bufev_f));
|
|
|
|
|
2009-04-13 03:08:11 +00:00
|
|
|
return downcast(bufev_f);
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
static void
|
|
|
|
be_filter_destruct(struct bufferevent *bev)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = upcast(bev);
|
|
|
|
assert(bevf);
|
|
|
|
if (bevf->free_context)
|
|
|
|
bevf->free_context(bevf->context);
|
|
|
|
|
2009-04-13 03:08:11 +00:00
|
|
|
if (bevf->bev.options & BEV_OPT_CLOSE_ON_FREE)
|
2009-02-02 19:22:13 +00:00
|
|
|
bufferevent_free(bevf->underlying);
|
2009-05-25 23:10:23 +00:00
|
|
|
|
|
|
|
_bufferevent_del_generic_timeout_cbs(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
static int
|
|
|
|
be_filter_enable(struct bufferevent *bev, short event)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = upcast(bev);
|
2009-05-25 23:10:23 +00:00
|
|
|
_bufferevent_generic_adj_timeouts(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
return bufferevent_enable(bevf->underlying, event);
|
|
|
|
}
|
|
|
|
|
|
|
|
static int
|
|
|
|
be_filter_disable(struct bufferevent *bev, short event)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = upcast(bev);
|
2009-05-25 23:10:23 +00:00
|
|
|
_bufferevent_generic_adj_timeouts(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
return bufferevent_disable(bevf->underlying, event);
|
|
|
|
}
|
|
|
|
|
|
|
|
static enum bufferevent_filter_result
|
|
|
|
be_filter_process_input(struct bufferevent_filtered *bevf,
|
|
|
|
enum bufferevent_flush_mode state,
|
|
|
|
int *processed_out)
|
|
|
|
{
|
|
|
|
enum bufferevent_filter_result res;
|
2009-04-13 03:08:11 +00:00
|
|
|
struct bufferevent *bev = downcast(bevf);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
if (state == BEV_NORMAL) {
|
|
|
|
/* If we're in 'normal' mode, don't urge data on the filter
|
|
|
|
* unless we're reading data and under our high-water mark.*/
|
2009-04-13 03:08:11 +00:00
|
|
|
if (!(bev->enabled & EV_READ) ||
|
2009-02-02 19:22:13 +00:00
|
|
|
be_readbuf_full(bevf, state))
|
|
|
|
return BEV_OK;
|
|
|
|
}
|
|
|
|
|
|
|
|
do {
|
2009-05-22 19:11:48 +00:00
|
|
|
ev_ssize_t limit = -1;
|
2009-04-13 03:08:11 +00:00
|
|
|
if (state == BEV_NORMAL && bev->wm_read.high)
|
|
|
|
limit = bev->wm_read.high -
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bev->input);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
res = bevf->process_in(bevf->underlying->input,
|
2009-04-13 03:08:11 +00:00
|
|
|
bev->input, limit, state, bevf->context);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
if (res == BEV_OK)
|
|
|
|
*processed_out = 1;
|
|
|
|
} while (res == BEV_OK &&
|
2009-04-13 03:08:11 +00:00
|
|
|
(bev->enabled & EV_READ) &&
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bevf->underlying->input) &&
|
2009-02-02 19:22:13 +00:00
|
|
|
!be_readbuf_full(bevf, state));
|
|
|
|
|
2009-05-25 23:10:23 +00:00
|
|
|
if (*processed_out)
|
|
|
|
BEV_RESET_GENERIC_READ_TIMEOUT(bev);
|
|
|
|
|
2009-02-02 19:22:13 +00:00
|
|
|
return res;
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
static enum bufferevent_filter_result
|
|
|
|
be_filter_process_output(struct bufferevent_filtered *bevf,
|
|
|
|
enum bufferevent_flush_mode state,
|
|
|
|
int *processed_out)
|
|
|
|
{
|
2009-07-17 17:46:17 +00:00
|
|
|
/* Requires references and lock: might call writecb */
|
2009-02-02 19:22:13 +00:00
|
|
|
enum bufferevent_filter_result res = BEV_OK;
|
|
|
|
struct bufferevent *bufev = downcast(bevf);
|
|
|
|
int again = 0;
|
|
|
|
|
|
|
|
if (state == BEV_NORMAL) {
|
|
|
|
/* If we're in 'normal' mode, don't urge data on the
|
|
|
|
* filter unless we're writing data, and the underlying
|
|
|
|
* bufferevent is accepting data, and we have data to
|
|
|
|
* give the filter. If we're in 'flush' or 'finish',
|
|
|
|
* call the filter no matter what. */
|
|
|
|
if (!(bufev->enabled & EV_WRITE) ||
|
|
|
|
be_underlying_writebuf_full(bevf, state) ||
|
2009-04-17 06:56:09 +00:00
|
|
|
!evbuffer_get_length(bufev->output))
|
2009-02-02 19:22:13 +00:00
|
|
|
return BEV_OK;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* disable the callback that calls this function
|
|
|
|
when the user adds to the output buffer. */
|
|
|
|
evbuffer_cb_set_flags(bufev->output, bevf->outbuf_cb, 0);
|
|
|
|
|
|
|
|
do {
|
|
|
|
int processed = 0;
|
|
|
|
again = 0;
|
|
|
|
|
|
|
|
do {
|
2009-05-22 19:11:48 +00:00
|
|
|
ev_ssize_t limit = -1;
|
2009-02-02 19:22:13 +00:00
|
|
|
if (state == BEV_NORMAL &&
|
|
|
|
bevf->underlying->wm_write.high)
|
|
|
|
limit = bevf->underlying->wm_write.high -
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bevf->underlying->output);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
2009-04-13 03:08:11 +00:00
|
|
|
res = bevf->process_out(downcast(bevf)->output,
|
2009-02-02 19:22:13 +00:00
|
|
|
bevf->underlying->output,
|
|
|
|
limit,
|
|
|
|
state,
|
|
|
|
bevf->context);
|
|
|
|
|
|
|
|
if (res == BEV_OK)
|
|
|
|
processed = *processed_out = 1;
|
|
|
|
} while (/* Stop if the filter wasn't successful...*/
|
|
|
|
res == BEV_OK &&
|
|
|
|
/* Or if we aren't writing any more. */
|
|
|
|
(bufev->enabled & EV_WRITE) &&
|
|
|
|
/* Of if we have nothing more to write and we are
|
|
|
|
* not flushing. */
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bufev->output) &&
|
2009-02-02 19:22:13 +00:00
|
|
|
/* Or if we have filled the underlying output buffer. */
|
|
|
|
!be_underlying_writebuf_full(bevf,state));
|
|
|
|
|
|
|
|
if (processed && bufev->writecb &&
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bufev->output) <= bufev->wm_write.low) {
|
2009-02-02 19:22:13 +00:00
|
|
|
/* call the write callback.*/
|
2009-04-17 23:12:34 +00:00
|
|
|
_bufferevent_run_writecb(bufev);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
|
|
|
if (res == BEV_OK &&
|
|
|
|
(bufev->enabled & EV_WRITE) &&
|
2009-04-17 06:56:09 +00:00
|
|
|
evbuffer_get_length(bufev->output) &&
|
2009-02-02 19:22:13 +00:00
|
|
|
!be_underlying_writebuf_full(bevf, state)) {
|
|
|
|
again = 1;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
} while (again);
|
|
|
|
|
|
|
|
/* reenable the outbuf_cb */
|
|
|
|
evbuffer_cb_set_flags(bufev->output,bevf->outbuf_cb,
|
|
|
|
EVBUFFER_CB_ENABLED);
|
|
|
|
|
2009-05-25 23:10:23 +00:00
|
|
|
if (*processed_out)
|
|
|
|
BEV_RESET_GENERIC_WRITE_TIMEOUT(bufev);
|
|
|
|
|
2009-02-02 19:22:13 +00:00
|
|
|
return res;
|
|
|
|
}
|
|
|
|
|
|
|
|
/* Called when the size of our outbuf changes. */
|
|
|
|
static void
|
|
|
|
bufferevent_filtered_outbuf_cb(struct evbuffer *buf,
|
2009-04-03 14:27:03 +00:00
|
|
|
const struct evbuffer_cb_info *cbinfo, void *arg)
|
2009-02-02 19:22:13 +00:00
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = arg;
|
2009-07-17 17:46:17 +00:00
|
|
|
struct bufferevent *bev = downcast(bevf);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
2009-04-03 14:27:03 +00:00
|
|
|
if (cbinfo->n_added) {
|
2009-02-02 19:22:13 +00:00
|
|
|
int processed_any = 0;
|
|
|
|
/* Somebody added more data to the output buffer. Try to
|
|
|
|
* process it, if we should. */
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_incref_and_lock(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
be_filter_process_output(bevf, BEV_NORMAL, &processed_any);
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_decref_and_unlock(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
/* Called when the underlying socket has read. */
|
|
|
|
static void
|
|
|
|
be_filter_readcb(struct bufferevent *underlying, void *_me)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = _me;
|
|
|
|
enum bufferevent_filter_result res;
|
|
|
|
enum bufferevent_flush_mode state;
|
|
|
|
struct bufferevent *bufev = downcast(bevf);
|
|
|
|
int processed_any = 0;
|
|
|
|
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_incref_and_lock(bufev);
|
|
|
|
|
2009-02-02 19:22:13 +00:00
|
|
|
if (bevf->got_eof)
|
|
|
|
state = BEV_FINISHED;
|
|
|
|
else
|
|
|
|
state = BEV_NORMAL;
|
|
|
|
|
|
|
|
res = be_filter_process_input(bevf, state, &processed_any);
|
|
|
|
|
2009-07-17 17:46:17 +00:00
|
|
|
/* XXX This should be in process_input, not here. There are
|
|
|
|
* other places that can call process-input, and they should
|
|
|
|
* force readcb calls as needed. */
|
2009-02-02 19:22:13 +00:00
|
|
|
if (processed_any &&
|
2009-05-25 23:10:23 +00:00
|
|
|
evbuffer_get_length(bufev->input) >= bufev->wm_read.low &&
|
|
|
|
bufev->readcb != NULL)
|
2009-04-17 23:12:34 +00:00
|
|
|
_bufferevent_run_readcb(bufev);
|
2009-07-17 17:46:17 +00:00
|
|
|
|
|
|
|
_bufferevent_decref_and_unlock(bufev);
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/* Called when the underlying socket has drained enough that we can write to
|
|
|
|
it. */
|
|
|
|
static void
|
|
|
|
be_filter_writecb(struct bufferevent *underlying, void *_me)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = _me;
|
2009-07-17 17:46:17 +00:00
|
|
|
struct bufferevent *bev = downcast(bevf);
|
2009-02-02 19:22:13 +00:00
|
|
|
int processed_any = 0;
|
|
|
|
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_incref_and_lock(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
be_filter_process_output(bevf, BEV_NORMAL, &processed_any);
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_decref_and_unlock(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/* Called when the underlying socket has given us an error */
|
|
|
|
static void
|
2009-05-25 23:11:20 +00:00
|
|
|
be_filter_eventcb(struct bufferevent *underlying, short what, void *_me)
|
2009-02-02 19:22:13 +00:00
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = _me;
|
2009-04-13 03:08:11 +00:00
|
|
|
struct bufferevent *bev = downcast(bevf);
|
2009-02-02 19:22:13 +00:00
|
|
|
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_incref_and_lock(bev);
|
2009-05-25 23:11:20 +00:00
|
|
|
/* All we can really to is tell our own eventcb. */
|
2009-04-13 03:08:11 +00:00
|
|
|
if (bev->errorcb)
|
2009-05-25 23:11:20 +00:00
|
|
|
_bufferevent_run_eventcb(bev, what);
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_decref_and_unlock(bev);
|
2009-02-02 19:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
static int
|
|
|
|
be_filter_flush(struct bufferevent *bufev,
|
|
|
|
short iotype, enum bufferevent_flush_mode mode)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf = upcast(bufev);
|
|
|
|
int processed_any = 0;
|
|
|
|
assert(bevf);
|
|
|
|
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_incref_and_lock(bufev);
|
|
|
|
|
2009-02-02 19:22:13 +00:00
|
|
|
if (iotype & EV_READ) {
|
|
|
|
be_filter_process_input(bevf, mode, &processed_any);
|
|
|
|
}
|
|
|
|
if (iotype & EV_WRITE) {
|
|
|
|
be_filter_process_output(bevf, mode, &processed_any);
|
|
|
|
}
|
|
|
|
/* XXX check the return value? */
|
|
|
|
/* XXX does this want to recursively call lower-level flushes? */
|
|
|
|
bufferevent_flush(bevf->underlying, iotype, mode);
|
|
|
|
|
2009-07-17 17:46:17 +00:00
|
|
|
_bufferevent_decref_and_unlock(bufev);
|
|
|
|
|
2009-02-02 19:22:13 +00:00
|
|
|
return processed_any;
|
|
|
|
}
|
2009-05-13 20:37:21 +00:00
|
|
|
|
|
|
|
static int
|
|
|
|
be_filter_ctrl(struct bufferevent *bev, enum bufferevent_ctrl_op op,
|
|
|
|
union bufferevent_ctrl_data *data)
|
|
|
|
{
|
|
|
|
struct bufferevent_filtered *bevf;
|
|
|
|
switch(op) {
|
|
|
|
case BEV_CTRL_GET_UNDERLYING:
|
|
|
|
bevf = upcast(bev);
|
|
|
|
data->ptr = bevf->underlying;
|
|
|
|
return 0;
|
|
|
|
case BEV_CTRL_GET_FD:
|
|
|
|
case BEV_CTRL_SET_FD:
|
|
|
|
default:
|
|
|
|
return -1;
|
|
|
|
}
|
|
|
|
}
|