FILTER_FORMAT_ERROR,
};
-/*!
- * Linked list of events.
- * Global events are appended to the list by append_event().
- * The usecount is the number of stored pointers to the element,
- * excluding the list pointers. So an element that is only in
- * the list has a usecount of 0, not 1.
- *
- * Clients have a pointer to the last event processed, and for each
- * of these clients we track the usecount of the elements.
- * If we have a pointer to an entry in the list, it is safe to navigate
- * it forward because elements will not be deleted, but only appended.
- * The worst that can happen is seeing the pointer still NULL.
+/*! \brief An AMI event message
*
- * When the usecount of an element drops to 0, and the element is the
- * first in the list, we can remove it. Removal is done within the
- * main thread, which is woken up for the purpose.
- *
- * For simplicity of implementation, we make sure the list is never empty.
+ * Events are created as reference counted immutable objects. Each
+ * active session that should receive the event is given a pointer to
+ * the event with a reference. Once the reference reaches 0 the event
+ * is destroyed. These events are stored in a vector on each session
+ * so that each session does not need to allocate a linked list
+ * wrapper for each event. Instead the vector grows as needed with
+ * excess each time to reduce the number of overall reallocations.
*/
struct eventqent {
- int usecount; /*!< # of clients who still need the event */
+ /*! \brief Category this event belongs to */
int category;
- unsigned int seq; /*!< sequence number */
- struct timeval tv; /*!< When event was allocated */
+ /*! \brief Hash of the name of the event */
int event_name_hash;
- AST_RWLIST_ENTRY(eventqent) eq_next;
- char eventdata[1]; /*!< really variable size, allocated by append_event() */
+ /*! \brief The contents of the AMI event message */
+ struct ast_str *message;
};
-static AST_RWLIST_HEAD_STATIC(all_events, eventqent);
-
static int displayconnects = 1;
static int allowmultiplelogin = 1;
static int timestampevents;
{{ "restart", "gracefully", NULL }},
};
+static struct eventqent *eventqent_alloc(const char *event_name, int category);
+
static void acl_change_stasis_cb(void *data, struct stasis_subscription *sub, struct stasis_message *message);
static void acl_change_stasis_subscribe(void)
struct ast_iostream *stream; /*!< AMI stream */
int inuse; /*!< number of HTTP sessions using this entry */
int needdestroy; /*!< Whether an HTTP session should be destroyed */
- pthread_t waiting_thread; /*!< Sleeping thread using this descriptor */
+ ast_mutex_t http_thread_lock; /*!< Lock for protecting the HTTP thread waiting on this session */
+ pthread_t http_thread; /*!< HTTP thread waiting on this session */
uint32_t managerid; /*!< Unique manager identifier, 0 for AMI sessions */
time_t sessionstart; /*!< Session start time */
struct timeval sessionstart_tv; /*!< Session start time */
struct ao2_container *includefilters; /*!< Manager event filters - include list */
struct ao2_container *excludefilters; /*!< Manager event filters - exclude list */
struct ast_variable *chanvars; /*!< Channel variables to set for originate */
- int send_events; /*!< XXX what ? */
- struct eventqent *last_ev; /*!< last event processed. */
+ int send_events; /*!< Event categories to send to this session */
int writetimeout; /*!< Timeout for ast_carefulwrite() */
time_t authstart;
- int pending_event; /*!< Pending events indicator in case when waiting_thread is NULL */
time_t noncetime; /*!< Timer for nonce value expiration */
unsigned long oldnonce; /*!< Stale nonce value */
unsigned long nc; /*!< incremental nonce counter */
unsigned int kicked:1; /*!< Flag set if session is forcibly kicked */
- ast_mutex_t notify_lock; /*!< Lock for notifying this session of events */
+ int alert_pipe[2]; /*!< Pipe for alerting this session */
+ AST_VECTOR(, struct eventqent *) pending_events; /*!< Queue of pending events to send */
+ ast_mutex_t pending_events_lock; /*!< Lock for pending events queue */
AST_LIST_HEAD_NOLOCK(mansession_datastores, ast_datastore) datastores; /*!< Data stores on the session */
AST_LIST_ENTRY(mansession_session) list;
};
return (webmanager_enabled && manager_enabled);
}
-/*!
- * Grab a reference to the last event, update usecount as needed.
- * Can handle a NULL pointer.
- */
-static struct eventqent *grab_last(void)
-{
- struct eventqent *ret;
-
- AST_RWLIST_WRLOCK(&all_events);
- ret = AST_RWLIST_LAST(&all_events);
- /* the list is never empty now, but may become so when
- * we optimize it in the future, so be prepared.
- */
- if (ret) {
- ast_atomic_fetchadd_int(&ret->usecount, 1);
- }
- AST_RWLIST_UNLOCK(&all_events);
- return ret;
-}
-
-/*!
- * Purge unused events. Remove elements from the head
- * as long as their usecount is 0 and there is a next element.
- */
-static void purge_events(void)
-{
- struct eventqent *ev;
- struct timeval now = ast_tvnow();
-
- AST_RWLIST_WRLOCK(&all_events);
- while ( (ev = AST_RWLIST_FIRST(&all_events)) &&
- ev->usecount == 0 && AST_RWLIST_NEXT(ev, eq_next)) {
- AST_RWLIST_REMOVE_HEAD(&all_events, eq_next);
- ast_free(ev);
- }
-
- AST_RWLIST_TRAVERSE_SAFE_BEGIN(&all_events, ev, eq_next) {
- /* Never release the last event */
- if (!AST_RWLIST_NEXT(ev, eq_next)) {
- break;
- }
-
- /* 2.5 times whatever the HTTP timeout is (maximum 2.5 hours) is the maximum time that we will definitely cache an event */
- if (ev->usecount == 0 && ast_tvdiff_sec(now, ev->tv) > (httptimeout > 3600 ? 3600 : httptimeout) * 2.5) {
- AST_RWLIST_REMOVE_CURRENT(eq_next);
- ast_free(ev);
- }
- }
- AST_RWLIST_TRAVERSE_SAFE_END;
- AST_RWLIST_UNLOCK(&all_events);
-}
-
/*!
* helper functions to convert back and forth between
* string and numeric representation of set of flags
ast_free(entry->string_filter);
}
+static void session_notify(struct mansession_session *session)
+{
+ if (!session->managerid) {
+ /* TCP based connections are easy and use an alert pipe, which requires no locking */
+ ast_alertpipe_write(session->alert_pipe);
+ } else {
+ /* HTTP based connections require locking as multiple HTTP connections can end up using the same session */
+ ast_mutex_lock(&session->http_thread_lock);
+ if (session->http_thread != AST_PTHREADT_NULL) {
+ pthread_kill(session->http_thread, SIGURG);
+ }
+ ast_mutex_unlock(&session->http_thread_lock);
+ }
+}
+
static void session_destructor(void *obj)
{
struct mansession_session *session = obj;
- struct eventqent *eqe = session->last_ev;
struct ast_datastore *datastore;
/* Get rid of each of the data stores on the session */
ast_datastore_free(datastore);
}
- if (eqe) {
- ast_atomic_fetchadd_int(&eqe->usecount, -1);
- }
if (session->chanvars) {
ast_variables_destroy(session->chanvars);
}
ao2_t_ref(session->excludefilters, -1, "decrement ref for exclude container, should be last one");
}
- ast_mutex_destroy(&session->notify_lock);
+ AST_VECTOR_CALLBACK_VOID(&session->pending_events, ao2_cleanup);
+ AST_VECTOR_FREE(&session->pending_events);
+ ast_mutex_destroy(&session->http_thread_lock);
+ ast_mutex_destroy(&session->pending_events_lock);
+
+ /* On initialization the alert pipe is set to -1, ensuring it is safe to call close */
+ ast_alertpipe_close(session->alert_pipe);
}
/*! \brief Allocate manager session structure and add it to the list of sessions */
-static struct mansession_session *build_mansession(const struct ast_sockaddr *addr)
+static struct mansession_session *build_mansession(const struct ast_sockaddr *addr, unsigned int needs_alert_pipe)
{
struct ao2_container *sessions;
struct mansession_session *newsession;
return NULL;
}
+ AST_VECTOR_INIT(&newsession->pending_events, 0);
+ ast_alertpipe_clear(newsession->alert_pipe);
+ ast_mutex_init(&newsession->pending_events_lock);
+ ast_mutex_init(&newsession->http_thread_lock);
+ AST_LIST_HEAD_INIT_NOLOCK(&newsession->datastores);
+
newsession->includefilters = ao2_container_alloc_list(AO2_ALLOC_OPT_LOCK_MUTEX, 0, NULL, NULL);
newsession->excludefilters = ao2_container_alloc_list(AO2_ALLOC_OPT_LOCK_MUTEX, 0, NULL, NULL);
if (!newsession->includefilters || !newsession->excludefilters) {
return NULL;
}
- newsession->waiting_thread = AST_PTHREADT_NULL;
+ newsession->http_thread = AST_PTHREADT_NULL;
newsession->writetimeout = 100;
- newsession->send_events = -1;
+ newsession->send_events = 0; /* Until authenticated this session will not receive events */
ast_sockaddr_copy(&newsession->addr, addr);
- ast_mutex_init(&newsession->notify_lock);
+ if (needs_alert_pipe && ast_alertpipe_init(newsession->alert_pipe)) {
+ ao2_ref(newsession, -1);
+ return NULL;
+ }
sessions = ao2_global_obj_ref(mgr_sessions);
if (sessions) {
fd = ast_iostream_get_fd(session->stream);
found = fd;
ast_cli(a->fd, "Kicking manager session connected using file descriptor %d\n", fd);
- ast_mutex_lock(&session->notify_lock);
session->kicked = 1;
- if (session->waiting_thread != AST_PTHREADT_NULL) {
- pthread_kill(session->waiting_thread, SIGURG);
- }
- ast_mutex_unlock(&session->notify_lock);
ao2_unlock(session);
+ session_notify(session);
unref_mansession(session);
break;
}
return CLI_SUCCESS;
}
-/*! \brief CLI command manager list eventq */
-/* Should change to "manager show connected" */
-static char *handle_showmaneventq(struct ast_cli_entry *e, int cmd, struct ast_cli_args *a)
-{
- struct eventqent *s;
- switch (cmd) {
- case CLI_INIT:
- e->command = "manager show eventq";
- e->usage =
- "Usage: manager show eventq\n"
- " Prints a listing of all events pending in the Asterisk manger\n"
- "event queue.\n";
- return NULL;
- case CLI_GENERATE:
- return NULL;
- }
- AST_RWLIST_RDLOCK(&all_events);
- AST_RWLIST_TRAVERSE(&all_events, s, eq_next) {
- ast_cli(a->fd, "Usecount: %d\n", s->usecount);
- ast_cli(a->fd, "Category: %d\n", s->category);
- ast_cli(a->fd, "Event:\n%s", s->eventdata);
- }
- AST_RWLIST_UNLOCK(&all_events);
-
- return CLI_SUCCESS;
-}
-
static int reload_module(void);
/*! \brief CLI command manager reload */
return CLI_SUCCESS;
}
-static struct eventqent *advance_event(struct eventqent *e)
-{
- struct eventqent *next;
-
- AST_RWLIST_RDLOCK(&all_events);
- if ((next = AST_RWLIST_NEXT(e, eq_next))) {
- ast_atomic_fetchadd_int(&next->usecount, 1);
- ast_atomic_fetchadd_int(&e->usecount, -1);
- }
- AST_RWLIST_UNLOCK(&all_events);
- return next;
-}
-
#define GET_HEADER_FIRST_MATCH 0
#define GET_HEADER_LAST_MATCH 1
#define GET_HEADER_SKIP_EMPTY 2
int maskint = strings_to_mask(eventmask);
ao2_lock(s->session);
+ /* Limit what event categories are sent to the client based on read permissions */
if (maskint >= 0) {
- s->session->send_events = maskint;
+ s->session->send_events = s->session->readperm & maskint;
}
ao2_unlock(s->session);
struct ast_manager_user *user = NULL;
regex_t *regex_filter;
struct ao2_iterator filter_iter;
+ const char *events;
if (ast_strlen_zero(username)) { /* missing username */
return -1;
s->session->sessionstart = time(NULL);
s->session->sessionstart_tv = ast_tvnow();
- set_eventmask(s, astman_get_header(m, "Events"));
+
+ /* When authenticating an Events header can be used to specify what event
+ * categories are desired. If no Events header is present or it is empty
+ * then the default read permissions of the user should be used.
+ */
+ events = astman_get_header(m, "Events");
+ if (!ast_strlen_zero(events)) {
+ set_eventmask(s, events);
+ } else {
+ s->session->send_events = s->session->readperm;
+ }
report_auth_success(s);
return 0;
}
-static int action_waitevent(struct mansession *s, const struct message *m)
+static int action_waitevent_tcp(struct mansession *s, const struct message *m, const char *id, int timeout)
{
- const char *timeouts = astman_get_header(m, "Timeout");
- int timeout = -1;
- int x;
- int needexit = 0;
- const char *id = astman_get_header(m, "ActionID");
- char idText[256];
+ struct pollfd pfds[2];
+ int x, res;
+ size_t events_count;
+ struct eventqent **events;
- if (!ast_strlen_zero(id)) {
- snprintf(idText, sizeof(idText), "ActionID: %s\r\n", id);
- } else {
- idText[0] = '\0';
+ /* Since we don't have access to the reusable pollfd, we set it up each time */
+ memset(pfds, 0, sizeof(pfds));
+ pfds[0].fd = ast_iostream_get_fd(s->session->stream);
+ pfds[0].events = POLLIN | POLLPRI;
+ pfds[1].fd = ast_alertpipe_readfd(s->session->alert_pipe);
+ pfds[1].events = POLLIN | POLLPRI;
+
+ ast_debug(1, "Starting waiting for an event!\n");
+
+ for (x = 0; x < timeout || timeout < 0; x++) {
+ /*
+ * Poll on both the TCP connection as well as the alert pipe for events,
+ * and when it comes to the timeout we get it passed in as seconds and the
+ * poll will wake up every second. The for loop then enforces the timeout.
+ */
+ res = ast_poll(pfds, 2, 1000);
+ if (res == -1) {
+ if (errno == EINTR || errno == EAGAIN) {
+ continue;
+ }
+ ast_log(LOG_WARNING, "poll() returned error: %s\n", strerror(errno));
+ return -1;
+ }
+
+ /* If the alert pipe has data that means there are events waiting */
+ if (pfds[1].revents) {
+ /* If the session has been kicked go no further */
+ if (s->session->kicked) {
+ ast_debug(1, "Manager session has been kicked\n");
+ return -1;
+ }
+
+ ast_alertpipe_read(s->session->alert_pipe);
+ break;
+ }
+
+ /* If any data came from the TCP connection break now or else poll will return immediately */
+ if (pfds[0].revents) {
+ break;
+ }
}
- if (!ast_strlen_zero(timeouts)) {
- sscanf(timeouts, "%30i", &timeout);
- if (timeout < -1) {
- timeout = -1;
+ ast_debug(1, "Finished waiting for an event!\n");
+
+ /* To reduce contention we lock only long enough to steal the events */
+ ast_mutex_lock(&s->session->pending_events_lock);
+ events_count = AST_VECTOR_SIZE(&s->session->pending_events);
+ events = AST_VECTOR_STEAL_ELEMENTS(&s->session->pending_events);
+ ast_mutex_unlock(&s->session->pending_events_lock);
+
+ ao2_lock(s->session);
+ astman_send_response(s, m, "Success", "Waiting for Event completed.");
+
+ for (x = 0; x < events_count; x++) {
+ struct eventqent *eqe = events[x];
+
+ if (((s->session->send_events & eqe->category) == eqe->category) &&
+ should_send_event(s->session->includefilters, s->session->excludefilters, eqe)) {
+ astman_append(s, "%s", ast_str_buffer(eqe->message));
}
- /* XXX maybe put an upper bound, or prevent the use of 0 ? */
+
+ ao2_ref(eqe, -1);
}
- ast_mutex_lock(&s->session->notify_lock);
- if (s->session->waiting_thread != AST_PTHREADT_NULL) {
- pthread_kill(s->session->waiting_thread, SIGURG);
+ astman_append(s,
+ "Event: WaitEventComplete\r\n"
+ "%s"
+ "\r\n", id);
+ ao2_unlock(s->session);
+
+ ast_free(events);
+
+ return 0;
+}
+
+static int action_waitevent_http(struct mansession *s, const struct message *m, const char *id, int timeout)
+{
+ time_t now;
+ int max, needexit = 0, x;
+
+ ast_mutex_lock(&s->session->http_thread_lock);
+ if (s->session->http_thread != AST_PTHREADT_NULL) {
+ pthread_kill(s->session->http_thread, SIGURG);
}
- ast_mutex_unlock(&s->session->notify_lock);
+ ast_mutex_unlock(&s->session->http_thread_lock);
ao2_lock(s->session);
- if (s->session->managerid) { /* AMI-over-HTTP session */
- /*
- * Make sure the timeout is within the expire time of the session,
- * as the client will likely abort the request if it does not see
- * data coming after some amount of time.
- */
- time_t now = time(NULL);
- int max = s->session->sessiontimeout - now - 10;
+ /*
+ * Make sure the timeout is within the expire time of the session,
+ * as the client will likely abort the request if it does not see
+ * data coming after some amount of time.
+ */
+ now = time(NULL);
+ max = s->session->sessiontimeout - now - 10;
- if (max < 0) { /* We are already late. Strange but possible. */
- max = 0;
- }
- if (timeout < 0 || timeout > max) {
- timeout = max;
- }
- if (!s->session->send_events) { /* make sure we record events */
- s->session->send_events = -1;
- }
+ if (max < 0) { /* We are already late. Strange but possible. */
+ max = 0;
+ }
+ if (timeout < 0 || timeout > max) {
+ timeout = max;
+ }
+ if (!s->session->send_events) { /* make sure we record events */
+ s->session->send_events = s->session->readperm;
}
ao2_unlock(s->session);
- ast_mutex_lock(&s->session->notify_lock);
- s->session->waiting_thread = pthread_self(); /* let new events wake up this thread */
- ast_mutex_unlock(&s->session->notify_lock);
+ ast_mutex_lock(&s->session->http_thread_lock);
+ s->session->http_thread = pthread_self(); /* let new events wake up this thread */
+ ast_mutex_unlock(&s->session->http_thread_lock);
ast_debug(1, "Starting waiting for an event!\n");
for (x = 0; x < timeout || timeout < 0; x++) {
- ao2_lock(s->session);
- if (AST_RWLIST_NEXT(s->session->last_ev, eq_next)) {
+ ast_mutex_lock(&s->session->pending_events_lock);
+ if (AST_VECTOR_SIZE(&s->session->pending_events)) {
needexit = 1;
}
+ ast_mutex_unlock(&s->session->pending_events_lock);
+
+ ao2_lock(s->session);
if (s->session->needdestroy) {
needexit = 1;
}
* The way we deal with it is not very nice: newcomers kick out the previous
* HTTP session. XXX this needs to be improved.
*/
- ast_mutex_lock(&s->session->notify_lock);
- if (s->session->waiting_thread != pthread_self()) {
+ ast_mutex_lock(&s->session->http_thread_lock);
+ if (s->session->http_thread != pthread_self()) {
needexit = 1;
}
- ast_mutex_unlock(&s->session->notify_lock);
+ ast_mutex_unlock(&s->session->http_thread_lock);
if (needexit) {
break;
}
- if (s->session->managerid == 0) { /* AMI session */
- if (ast_wait_for_input(ast_iostream_get_fd(s->session->stream), 1000)) {
- break;
- }
- } else { /* HTTP session */
- sleep(1);
- }
+
+ sleep(1);
}
ast_debug(1, "Finished waiting for an event!\n");
- ast_mutex_lock(&s->session->notify_lock);
- if (s->session->waiting_thread == pthread_self()) {
- struct eventqent *eqe = s->session->last_ev;
+ ast_mutex_lock(&s->session->http_thread_lock);
+ if (s->session->http_thread == pthread_self()) {
+ size_t events_count;
+ struct eventqent **events;
- s->session->waiting_thread = AST_PTHREADT_NULL;
- ast_mutex_unlock(&s->session->notify_lock);
+ s->session->http_thread = AST_PTHREADT_NULL;
+ ast_mutex_unlock(&s->session->http_thread_lock);
+
+ /* To reduce contention we lock only long enough to steal the events */
+ ast_mutex_lock(&s->session->pending_events_lock);
+ events_count = AST_VECTOR_SIZE(&s->session->pending_events);
+ events = AST_VECTOR_STEAL_ELEMENTS(&s->session->pending_events);
+ ast_mutex_unlock(&s->session->pending_events_lock);
ao2_lock(s->session);
astman_send_response(s, m, "Success", "Waiting for Event completed.");
- while ((eqe = advance_event(eqe))) {
- if (((s->session->readperm & eqe->category) == eqe->category)
- && ((s->session->send_events & eqe->category) == eqe->category)
- && should_send_event(s->session->includefilters, s->session->excludefilters, eqe)) {
- astman_append(s, "%s", eqe->eventdata);
+ for (x = 0; x < events_count; x++) {
+ struct eventqent *eqe = events[x];
+
+ if (((s->session->send_events & eqe->category) == eqe->category) &&
+ should_send_event(s->session->includefilters, s->session->excludefilters, eqe)) {
+ astman_append(s, "%s", ast_str_buffer(eqe->message));
}
- s->session->last_ev = eqe;
+
+ ao2_ref(eqe, -1);
}
astman_append(s,
"Event: WaitEventComplete\r\n"
"%s"
- "\r\n", idText);
+ "\r\n", id);
ao2_unlock(s->session);
+ ast_free(events);
} else {
- ast_mutex_unlock(&s->session->notify_lock);
+ ast_mutex_unlock(&s->session->http_thread_lock);
ast_debug(1, "Abandoning event request!\n");
}
return 0;
}
+static int action_waitevent(struct mansession *s, const struct message *m)
+{
+ const char *timeouts = astman_get_header(m, "Timeout");
+ int timeout = -1;
+ const char *id = astman_get_header(m, "ActionID");
+ char idText[256];
+
+ if (!ast_strlen_zero(id)) {
+ snprintf(idText, sizeof(idText), "ActionID: %s\r\n", id);
+ } else {
+ idText[0] = '\0';
+ }
+
+ if (!ast_strlen_zero(timeouts)) {
+ sscanf(timeouts, "%30i", &timeout);
+ if (timeout < -1) {
+ timeout = -1;
+ }
+ /* XXX maybe put an upper bound, or prevent the use of 0 ? */
+ }
+
+ if (!s->session->managerid) {
+ return action_waitevent_tcp(s, m, idText, timeout);
+ } else {
+ return action_waitevent_http(s, m, idText, timeout);
+ }
+}
+
static int action_listcommands(struct mansession *s, const struct message *m)
{
struct manager_action *cur;
}
astman_send_ack(s, m, "Authentication accepted");
if ((s->session->send_events & EVENT_FLAG_SYSTEM)
- && (s->session->readperm & EVENT_FLAG_SYSTEM)
&& ast_fully_booted) {
struct ast_str *auth = ast_str_alloca(MAX_AUTH_PERM_STRING);
const char *cat_str = authority_to_str(EVENT_FLAG_SYSTEM, &auth);
/* We're looking at the entire event data */
if (!filter_entry->header_name) {
- match = match_eventdata(filter_entry, eqe->eventdata);
+ match = match_eventdata(filter_entry, ast_str_buffer(eqe->message));
goto done;
}
/* We're looking at a specific header */
- line_buffer_start = ast_strdup(eqe->eventdata);
+ line_buffer_start = ast_strdup(ast_str_buffer(eqe->message));
line_buffer = line_buffer_start;
if (!line_buffer_start) {
goto done;
int result = 0;
if (manager_debug) {
- ast_verbose("<-- Examining AMI event (%u): -->\n%s\n", eqe->event_name_hash, eqe->eventdata);
+ ast_verbose("<-- Examining AMI event (%u): -->\n%s\n", eqe->event_name_hash, ast_str_buffer(eqe->message));
} else {
- ast_debug(4, "Examining AMI event (%u):\n%s\n", eqe->event_name_hash, eqe->eventdata);
+ ast_debug(4, "Examining AMI event (%u):\n%s\n", eqe->event_name_hash, ast_str_buffer(eqe->message));
}
if (!ao2_container_count(includefilters) && !ao2_container_count(excludefilters)) {
return 1; /* no filtering means match all */
/*
* This is a basic test of whether a single event matches a single filter.
*/
- eqe = ast_calloc(1, sizeof(*eqe) + strlen(parsing_filter_tests[i].test_event_payload) + 1);
+ eqe = eventqent_alloc(parsing_filter_tests[i].test_event_name, 0);
if (!eqe) {
ast_test_status_update(test, "Failed to allocate eventqent\n");
res = AST_TEST_FAIL;
ao2_ref(filter_entry, -1);
break;
}
- strcpy(eqe->eventdata, parsing_filter_tests[i].test_event_payload); /* Safe */
- eqe->event_name_hash = ast_str_hash(parsing_filter_tests[i].test_event_name);
+ ast_str_set(&eqe->message, 0, "%s", parsing_filter_tests[i].test_event_payload);
send_event = should_send_event(includefilters, excludefilters, eqe);
if (send_event != parsing_filter_tests[i].expected_should_send_event) {
char *escaped = ast_escape_c_alloc(parsing_filter_tests[i].test_event_payload);
res = AST_TEST_FAIL;
}
loop_cleanup:
- ast_free(eqe);
+ ao2_cleanup(eqe);
ao2_cleanup(filter_entry);
}
int send_event = 0;
struct eventqent *eqe = NULL;
- eqe = ast_calloc(1, sizeof(*eqe) + strlen(events_for_matching[i].payload) + 1);
+ eqe = eventqent_alloc(events_for_matching[i].event_name, 0);
if (!eqe) {
ast_test_status_update(test, "Failed to allocate eventqent\n");
res = AST_TEST_FAIL;
break;
}
- strcpy(eqe->eventdata, events_for_matching[i].payload); /* Safe */
- eqe->event_name_hash = ast_str_hash(events_for_matching[i].event_name);
+ ast_str_set(&eqe->message, 0, "%s", events_for_matching[i].payload);
send_event = should_send_event(includefilters, excludefilters, eqe);
if (send_event != events_for_matching[i].expected_should_send_event) {
char *escaped = ast_escape_c_alloc(events_for_matching[i].payload);
ast_free(escaped);
res = AST_TEST_FAIL;
}
- ast_free(eqe);
+ ao2_ref(eqe, -1);
}
ast_test_debug(test, "Tested %d events\n", i);
#endif
/*!
- * Send any applicable events to the client listening on this socket.
- * Wait only for a finite time on each event, and drop all events whether
- * they are successfully sent or not.
+ * \brief Send any pending events to the client listening on this socket.
+ *
+ * \param s The manager session to send events for.
*/
static int process_events(struct mansession *s)
{
int ret = 0;
+ size_t events_count;
+ struct eventqent **events;
+ int i;
ao2_lock(s->session);
- if (s->session->stream != NULL) {
- struct eventqent *eqe = s->session->last_ev;
- while ((eqe = advance_event(eqe))) {
- if (eqe->category == EVENT_FLAG_SHUTDOWN) {
- ast_debug(3, "Received CloseSession event\n");
- ret = -1;
- }
- if (!ret && s->session->authenticated &&
- (s->session->readperm & eqe->category) == eqe->category &&
- (s->session->send_events & eqe->category) == eqe->category) {
- if (should_send_event(s->session->includefilters, s->session->excludefilters, eqe)) {
- if (send_string(s, eqe->eventdata) < 0 || s->write_error)
- ret = -1; /* don't send more */
- }
- }
- s->session->last_ev = eqe;
+ if (!s->session->stream) {
+ ao2_unlock(s->session);
+ return 0;
+ }
+
+ /* To reduce contention we lock only long enough to steal the events */
+ ast_mutex_lock(&s->session->pending_events_lock);
+ events_count = AST_VECTOR_SIZE(&s->session->pending_events);
+ events = AST_VECTOR_STEAL_ELEMENTS(&s->session->pending_events);
+ ast_mutex_unlock(&s->session->pending_events_lock);
+
+ for (i = 0; i < events_count; i++) {
+ struct eventqent *eqe = events[i];
+
+ if (eqe->category == EVENT_FLAG_SHUTDOWN) {
+ ast_debug(3, "Received CloseSession event\n");
+ ret = -1;
+ /* We purposely don't break here so that we drop the reference on remaining events */
+ }
+
+ if (!ret && ((s->session->send_events & eqe->category) == eqe->category) &&
+ should_send_event(s->session->includefilters, s->session->excludefilters, eqe) &&
+ (send_string(s, ast_str_buffer(eqe->message)) < 0 || s->write_error)) {
+ ret = -1; /* don't send more */
}
+
+ ao2_ref(eqe, -1);
}
+
ao2_unlock(s->session);
+
+ ast_free(events);
+
return ret;
}
* Also note that we assume output to have at least "maxlen" space.
* \endverbatim
*/
-static int get_input(struct mansession *s, char *output)
+static int get_input(struct mansession *s, struct pollfd *pfds, char *output)
{
- int res, x;
+ int res = 0, x;
int maxlen = sizeof(s->session->inbuf) - 1;
char *src = s->session->inbuf;
int timeout = -1;
}
}
- ast_mutex_lock(&s->session->notify_lock);
- if (s->session->pending_event) {
- s->session->pending_event = 0;
- ast_mutex_unlock(&s->session->notify_lock);
- return 0;
- }
- s->session->waiting_thread = pthread_self();
- ast_mutex_unlock(&s->session->notify_lock);
+ res = ast_poll(pfds, 2, timeout);
- res = ast_wait_for_input(ast_iostream_get_fd(s->session->stream), timeout);
+ /* If polling resulted in an error, break out of the loop */
+ if (res < 0) {
+ /* If in reality we were interrupted, iterate again */
+ if (errno == EINTR || errno == EAGAIN) {
+ continue;
+ }
- ast_mutex_lock(&s->session->notify_lock);
- s->session->waiting_thread = AST_PTHREADT_NULL;
- ast_mutex_unlock(&s->session->notify_lock);
- }
- if (res < 0) {
- if (s->session->kicked) {
- ast_debug(1, "Manager session has been kicked\n");
+ ast_log(LOG_WARNING, "poll() returned error: %s\n", strerror(errno));
return -1;
+ } else if (!res) {
+ /* If nothing happened we loop again */
+ continue;
}
- /* If we get a signal from some other thread (typically because
- * there are new events queued), return 0 to notify the caller.
- */
- if (errno == EINTR || errno == EAGAIN) {
- return 0;
+
+ if (pfds[1].revents) {
+ /* There are pending events to send, or this session has been kicked */
+ if (s->session->kicked) {
+ ast_debug(1, "Manager session has been kicked\n");
+ return -1;
+ }
+ ast_alertpipe_read(s->session->alert_pipe);
+
+ if (process_events(s)) {
+ return -1;
+ }
}
- ast_log(LOG_WARNING, "poll() returned error: %s\n", strerror(errno));
- return -1;
- }
- ao2_lock(s->session);
- res = ast_iostream_read(s->session->stream, src + s->session->inlen, maxlen - s->session->inlen);
- if (res < 1) {
- res = -1; /* error return */
- } else {
- s->session->inlen += res;
- src[s->session->inlen] = '\0';
- res = 0;
+ if (pfds[0].revents) {
+ ao2_lock(s->session);
+ res = ast_iostream_read(s->session->stream, src + s->session->inlen, maxlen - s->session->inlen);
+ if (res < 1) {
+ res = -1; /* error return */
+ } else {
+ s->session->inlen += res;
+ src[s->session->inlen] = '\0';
+ res = 0;
+ }
+ ao2_unlock(s->session);
+ return res;
+ }
}
- ao2_unlock(s->session);
- return res;
+
+ return 0;
}
/*!
int res;
int hdr_loss;
time_t now;
+ struct pollfd pfds[2];
+
+ /* TCP based manager never has the underlying stream change, so we can setup the pollfd once */
+ memset(pfds, 0, sizeof(pfds));
+ pfds[0].fd = ast_iostream_get_fd(s->session->stream);
+ pfds[0].events = POLLIN | POLLPRI;
+ pfds[1].fd = ast_alertpipe_readfd(s->session->alert_pipe);
+ pfds[1].events = POLLIN | POLLPRI;
hdr_loss = 0;
for (;;) {
- /* Check if any events are pending and do them if needed */
- if (process_events(s)) {
- res = -1;
- break;
- }
- res = get_input(s, header_buf);
+ res = get_input(s, pfds, header_buf);
if (res == 0) {
/* No input line received. */
if (!s->session->authenticated) {
};
int res;
int arg = 1;
- struct ast_sockaddr ser_remote_address_tmp;
+ /* If this session would exceed the auth limit, reject it immediately */
if (ast_atomic_fetchadd_int(&unauth_sessions, +1) >= authlimit) {
ast_atomic_fetchadd_int(&unauth_sessions, -1);
goto done;
}
- ast_sockaddr_copy(&ser_remote_address_tmp, &ser->remote_address);
- session = build_mansession(&ser_remote_address_tmp);
-
- if (session == NULL) {
+ session = build_mansession(&ser->remote_address, 1);
+ if (!session) {
ast_atomic_fetchadd_int(&unauth_sessions, -1);
goto done;
}
ast_iostream_nonblock(ser->stream);
ao2_lock(session);
- /* Hook to the tail of the event queue */
- session->last_ev = grab_last();
ast_mutex_init(&s.lock);
/* these fields duplicate those in the 'ser' structure */
session->stream = s.stream = ser->stream;
- ast_sockaddr_copy(&session->addr, &ser_remote_address_tmp);
+ ast_sockaddr_copy(&session->addr, &ser->remote_address);
s.session = session;
- AST_LIST_HEAD_INIT_NOLOCK(&session->datastores);
-
if(time(&session->authstart) == -1) {
ast_log(LOG_ERROR, "error executing time(): %s; disconnecting client\n", strerror(errno));
ast_atomic_fetchadd_int(&unauth_sessions, -1);
return purged;
}
-/*! \brief
- * events are appended to a queue from where they
- * can be dispatched to clients.
+/*!
+ * \brief Destructor for manager event message
+ *
+ * \param obj The event message to free
*/
-static int append_event(const char *str, int event_name_hash, int category)
+static void eventqent_destructor(void *obj)
{
- struct eventqent *tmp = ast_malloc(sizeof(*tmp) + strlen(str));
- static int seq; /* sequence number */
+ struct eventqent *event = obj;
- if (!tmp) {
- return -1;
+ ast_free(event->message);
+}
+
+/*! \brief Initial size of the manager event message buffer */
+#define MANAGER_EVENT_BUF_INITSIZE 256
+
+/*!
+ * \brief Allocate a manager event
+ *
+ * \param event_name The name of the event
+ * \param category The category of the event
+ *
+ * \return non-NULL on success, NULL on failure
+ */
+static struct eventqent *eventqent_alloc(const char *event_name, int category)
+{
+ struct eventqent *event;
+
+ event = ao2_alloc_options(sizeof(*event), eventqent_destructor, AO2_ALLOC_OPT_LOCK_NOLOCK);
+ if (!event) {
+ return NULL;
}
- /* need to init all fields, because ast_malloc() does not */
- tmp->usecount = 0;
- tmp->category = category;
- tmp->seq = ast_atomic_fetchadd_int(&seq, 1);
- tmp->tv = ast_tvnow();
- tmp->event_name_hash = event_name_hash;
- AST_RWLIST_NEXT(tmp, eq_next) = NULL;
- strcpy(tmp->eventdata, str);
+ event->category = category;
+ event->event_name_hash = ast_str_hash(event_name);
- AST_RWLIST_WRLOCK(&all_events);
- AST_RWLIST_INSERT_TAIL(&all_events, tmp, eq_next);
- AST_RWLIST_UNLOCK(&all_events);
+ /*
+ * This event will most likely be queued to at least one session so we do message
+ * creation directly inside the event to avoid having to duplicate the string.
+ */
+ event->message = ast_str_create(MANAGER_EVENT_BUF_INITSIZE);
+ if (!event->message) {
+ ao2_ref(event, -1);
+ return NULL;
+ }
- return 0;
+ return event;
}
static void append_channel_vars(struct ast_str **pbuf, struct ast_channel *chan)
ao2_ref(vars, -1);
}
-/* XXX see if can be moved inside the function */
-AST_THREADSTORAGE(manager_event_buf);
-#define MANAGER_EVENT_BUF_INITSIZE 256
-
static int __attribute__((format(printf, 9, 0))) __manager_event_sessions_va(
struct ao2_container *sessions,
int category,
- const char *event,
+ const char *event_name,
int chancount,
struct ast_channel **chans,
const char *file,
const char *fmt,
va_list ap)
{
+ struct eventqent *event;
struct ast_str *auth = ast_str_alloca(MAX_AUTH_PERM_STRING);
const char *cat_str;
- struct timeval now;
- struct ast_str *buf;
int i;
- int event_name_hash;
if (!ast_strlen_zero(manager_disabledevents)) {
- if (ast_in_delimited_string(event, manager_disabledevents, ',')) {
- ast_debug(3, "AMI Event '%s' is globally disabled, skipping\n", event);
+ if (ast_in_delimited_string(event_name, manager_disabledevents, ',')) {
+ ast_debug(3, "AMI Event '%s' is globally disabled, skipping\n", event_name);
/* Event is globally disabled */
return -1;
}
}
- buf = ast_str_thread_get(&manager_event_buf, MANAGER_EVENT_BUF_INITSIZE);
- if (!buf) {
+ event = eventqent_alloc(event_name, category);
+ if (!event) {
return -1;
}
cat_str = authority_to_str(category, &auth);
- ast_str_set(&buf, 0,
+ ast_str_set(&event->message, 0,
"Event: %s\r\n"
"Privilege: %s\r\n",
- event, cat_str);
+ event_name, cat_str);
if (timestampevents) {
+ struct timeval now;
+
now = ast_tvnow();
- ast_str_append(&buf, 0,
+ ast_str_append(&event->message, 0,
"Timestamp: %ld.%06lu\r\n",
(long)now.tv_sec, (unsigned long) now.tv_usec);
}
if (manager_debug) {
static int seq;
- ast_str_append(&buf, 0,
+ ast_str_append(&event->message, 0,
"SequenceNumber: %d\r\n",
ast_atomic_fetchadd_int(&seq, 1));
- ast_str_append(&buf, 0,
+ ast_str_append(&event->message, 0,
"File: %s\r\n"
"Line: %d\r\n"
"Func: %s\r\n",
file, line, func);
}
if (!ast_strlen_zero(ast_config_AST_SYSTEM_NAME)) {
- ast_str_append(&buf, 0,
+ ast_str_append(&event->message, 0,
"SystemName: %s\r\n",
ast_config_AST_SYSTEM_NAME);
}
- ast_str_append_va(&buf, 0, fmt, ap);
+ ast_str_append_va(&event->message, 0, fmt, ap);
for (i = 0; i < chancount; i++) {
- append_channel_vars(&buf, chans[i]);
+ append_channel_vars(&event->message, chans[i]);
}
- ast_str_append(&buf, 0, "\r\n");
-
- event_name_hash = ast_str_hash(event);
-
- append_event(ast_str_buffer(buf), event_name_hash, category);
+ ast_str_append(&event->message, 0, "\r\n");
/* Wake up any sleeping sessions */
if (sessions) {
struct ao2_iterator iter;
struct mansession_session *session;
+ /*
+ * This purposely does not hold the session lock while checking send_events, as it
+ * could block for a long period of time. We also do a second check in the consumer
+ * of the pending_events queue such that if an event is added to the queue which should no
+ * longer be sent to the session it won't be.
+ */
iter = ao2_iterator_init(sessions, 0);
while ((session = ao2_iterator_next(&iter))) {
- ast_mutex_lock(&session->notify_lock);
- if (session->waiting_thread != AST_PTHREADT_NULL) {
- pthread_kill(session->waiting_thread, SIGURG);
- } else {
- /* We have an event to process, but the mansession is
- * not waiting for it. We still need to indicate that there
- * is an event waiting so that get_input processes the pending
- * event instead of polling.
- */
- session->pending_event = 1;
+ if (category == EVENT_FLAG_SHUTDOWN ||
+ (session->send_events & category) == category) {
+ ast_mutex_lock(&session->pending_events_lock);
+ if (!AST_VECTOR_APPEND(&session->pending_events, ao2_bump(event))) {
+ if (AST_VECTOR_SIZE(&session->pending_events) == 1) {
+ /* If this is the first event, wake up the session */
+ session_notify(session);
+ }
+ } else {
+ ao2_ref(event, -1);
+ }
+ ast_mutex_unlock(&session->pending_events_lock);
}
- ast_mutex_unlock(&session->notify_lock);
unref_mansession(session);
}
ao2_iterator_destroy(&iter);
AST_RWLIST_RDLOCK(&manager_hooks);
AST_RWLIST_TRAVERSE(&manager_hooks, hook, list) {
- hook->helper(category, event, ast_str_buffer(buf));
+ hook->helper(category, event_name, ast_str_buffer(event->message));
}
AST_RWLIST_UNLOCK(&manager_hooks);
}
+ ao2_ref(event, -1);
+
return 0;
}
/* Create new session.
* While it is not in the list we don't need any locking
*/
- if (!(session = build_mansession(remote_address))) {
+ if (!(session = build_mansession(remote_address, 0))) {
ast_http_request_close_on_completion(ser);
ast_http_error(ser, 500, "Server Error", "Internal Server Error (out of memory)");
return 0;
}
ao2_lock(session);
- session->send_events = 0;
session->inuse = 1;
/*!
* \note There is approximately a 1 in 1.8E19 chance that the following
*/
while ((session->managerid = ast_random() ^ (unsigned long) session) == 0) {
}
- session->last_ev = grab_last();
- AST_LIST_HEAD_INIT_NOLOCK(&session->datastores);
}
ao2_unlock(session);
blastaway = 1;
} else {
ast_debug(1, "Need destroy, but can't do it yet!\n");
- ast_mutex_lock(&session->notify_lock);
- if (session->waiting_thread != AST_PTHREADT_NULL) {
- pthread_kill(session->waiting_thread, SIGURG);
+ ast_mutex_lock(&session->http_thread_lock);
+ if (session->http_thread != AST_PTHREADT_NULL) {
+ pthread_kill(session->http_thread, SIGURG);
}
- ast_mutex_unlock(&session->notify_lock);
+ ast_mutex_unlock(&session->http_thread_lock);
session->inuse--;
}
} else {
* Create new session.
* While it is not in the list we don't need any locking
*/
- if (!(session = build_mansession(remote_address))) {
+ if (!(session = build_mansession(remote_address, 0))) {
ast_http_request_close_on_completion(ser);
ast_http_error(ser, 500, "Server Error", "Internal Server Error (out of memory)");
return 0;
ast_copy_string(session->username, u_username, sizeof(session->username));
session->managerid = nonce;
- session->last_ev = grab_last();
- AST_LIST_HEAD_INIT_NOLOCK(&session->datastores);
session->readperm = u_readperm;
session->writeperm = u_writeperm;
} else {
ser->poll_timeout = 5000;
}
- purge_events();
}
static struct ast_tls_config ami_tls_cfg;
AST_CLI_DEFINE(handle_showmancmds, "List manager interface commands"),
AST_CLI_DEFINE(handle_showmanconn, "List connected manager interface users"),
AST_CLI_DEFINE(handle_kickmanconn, "Kick a connected manager interface connection"),
- AST_CLI_DEFINE(handle_showmaneventq, "List manager interface queued events"),
AST_CLI_DEFINE(handle_showmanagers, "List configured manager users"),
AST_CLI_DEFINE(handle_showmanager, "Display information on a specific manager user"),
AST_CLI_DEFINE(handle_mandebug, "Show, enable, disable debugging of the manager code"),
__ast_custom_function_register(&managerclient_function, NULL);
ast_extension_state_add(NULL, NULL, manager_state_cb, NULL);
- /* Append placeholder event so master_eventq never runs dry */
- if (append_event("Event: Placeholder\r\n\r\n",
- ast_str_hash("Placeholder"), 0)) {
- return -1;
- }
-
#ifdef AST_XML_DOCS
temp_event_docs = ast_xmldoc_build_documentation("managerEvent");
if (temp_event_docs) {