diff --git a/main/manager.c b/main/manager.c index be968534a5..6cf07f55fa 100644 --- a/main/manager.c +++ b/main/manager.c @@ -131,37 +131,25 @@ enum add_filter_result { 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. +/*! \brief An AMI event message * - * 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. - * - * 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; @@ -230,6 +218,8 @@ static const struct { {{ "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) @@ -285,7 +275,8 @@ struct mansession_session { 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 */ @@ -300,16 +291,16 @@ struct mansession_session { 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; }; @@ -703,58 +694,6 @@ int ast_webmanager_check_enabled(void) 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 @@ -948,10 +887,24 @@ static void event_filter_destructor(void *obj) 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 */ @@ -960,9 +913,6 @@ static void session_destructor(void *obj) ast_datastore_free(datastore); } - if (eqe) { - ast_atomic_fetchadd_int(&eqe->usecount, -1); - } if (session->chanvars) { ast_variables_destroy(session->chanvars); } @@ -975,11 +925,17 @@ static void session_destructor(void *obj) 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; @@ -989,6 +945,12 @@ static struct mansession_session *build_mansession(const struct ast_sockaddr *ad 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) { @@ -996,12 +958,15 @@ static struct mansession_session *build_mansession(const struct ast_sockaddr *ad 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) { @@ -1473,13 +1438,9 @@ static char *handle_kickmanconn(struct ast_cli_entry *e, int cmd, struct ast_cli 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; } @@ -1546,33 +1507,6 @@ static char *handle_showmanconn(struct ast_cli_entry *e, int cmd, struct ast_cli 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 */ @@ -1595,19 +1529,6 @@ static char *handle_manager_reload(struct ast_cli_entry *e, int cmd, struct ast_ 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 @@ -2113,8 +2034,9 @@ static int set_eventmask(struct mansession *s, const char *eventmask) 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); @@ -2355,6 +2277,7 @@ static int authenticate(struct mansession *s, const struct message *m) 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; @@ -2436,7 +2359,17 @@ static int authenticate(struct mansession *s, const struct message *m) 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); @@ -3186,12 +3119,198 @@ static int action_createconfig(struct mansession *s, const struct message *m) return 0; } +static int action_waitevent_tcp(struct mansession *s, const struct message *m, const char *id, int timeout) +{ + struct pollfd pfds[2]; + int x, res; + size_t events_count; + struct eventqent **events; + + /* 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; + } + } + + 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)); + } + + ao2_ref(eqe, -1); + } + + 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->http_thread_lock); + + ao2_lock(s->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. + */ + 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 = s->session->readperm; + } + ao2_unlock(s->session); + + 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++) { + 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; + } + ao2_unlock(s->session); + /* We can have multiple HTTP session point to the same mansession entry. + * 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->http_thread_lock); + if (s->session->http_thread != pthread_self()) { + needexit = 1; + } + ast_mutex_unlock(&s->session->http_thread_lock); + if (needexit) { + break; + } + + sleep(1); + } + ast_debug(1, "Finished waiting for an event!\n"); + + ast_mutex_lock(&s->session->http_thread_lock); + if (s->session->http_thread == pthread_self()) { + size_t events_count; + struct eventqent **events; + + 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."); + 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)); + } + + ao2_ref(eqe, -1); + } + astman_append(s, + "Event: WaitEventComplete\r\n" + "%s" + "\r\n", id); + ao2_unlock(s->session); + ast_free(events); + } else { + 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; - int x; - int needexit = 0; const char *id = astman_get_header(m, "ActionID"); char idText[256]; @@ -3209,99 +3328,11 @@ static int action_waitevent(struct mansession *s, const struct message *m) /* XXX maybe put an upper bound, or prevent the use of 0 ? */ } - ast_mutex_lock(&s->session->notify_lock); - if (s->session->waiting_thread != AST_PTHREADT_NULL) { - pthread_kill(s->session->waiting_thread, SIGURG); - } - ast_mutex_unlock(&s->session->notify_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; - - 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; - } - } - 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_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)) { - needexit = 1; - } - if (s->session->needdestroy) { - needexit = 1; - } - ao2_unlock(s->session); - /* We can have multiple HTTP session point to the same mansession entry. - * 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()) { - needexit = 1; - } - ast_mutex_unlock(&s->session->notify_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); - } - } - 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; - - s->session->waiting_thread = AST_PTHREADT_NULL; - ast_mutex_unlock(&s->session->notify_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); - } - s->session->last_ev = eqe; - } - astman_append(s, - "Event: WaitEventComplete\r\n" - "%s" - "\r\n", idText); - ao2_unlock(s->session); + if (!s->session->managerid) { + return action_waitevent_tcp(s, m, idText, timeout); } else { - ast_mutex_unlock(&s->session->notify_lock); - ast_debug(1, "Abandoning event request!\n"); + return action_waitevent_http(s, m, idText, timeout); } - - return 0; } static int action_listcommands(struct mansession *s, const struct message *m) @@ -3394,7 +3425,6 @@ static int action_login(struct mansession *s, const struct message *m) } 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); @@ -5632,12 +5662,12 @@ static int filter_cmp_fn(void *obj, void *arg, void *data, int flags) /* 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; @@ -5673,9 +5703,9 @@ static int should_send_event(struct ao2_container *includefilters, 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 */ @@ -6272,15 +6302,14 @@ AST_TEST_DEFINE(eventfilter_test_creation) /* * 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); @@ -6292,7 +6321,7 @@ AST_TEST_DEFINE(eventfilter_test_creation) res = AST_TEST_FAIL; } loop_cleanup: - ast_free(eqe); + ao2_cleanup(eqe); ao2_cleanup(filter_entry); } @@ -6401,14 +6430,13 @@ AST_TEST_DEFINE(eventfilter_test_matching) 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); @@ -6419,7 +6447,7 @@ AST_TEST_DEFINE(eventfilter_test_matching) ast_free(escaped); res = AST_TEST_FAIL; } - ast_free(eqe); + ao2_ref(eqe, -1); } ast_test_debug(test, "Tested %d events\n", i); @@ -6428,35 +6456,52 @@ AST_TEST_DEFINE(eventfilter_test_matching) #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; } @@ -7173,9 +7218,9 @@ static int process_message(struct mansession *s, const struct message *m) * 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; @@ -7225,47 +7270,51 @@ static int get_input(struct mansession *s, char *output) } } - 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; + } + } + + 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; } - 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; - } - ao2_unlock(s->session); - return res; + return 0; } /*! @@ -7301,15 +7350,18 @@ static int do_message(struct mansession *s) 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) { @@ -7388,17 +7440,15 @@ static void *session_do(void *data) }; 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; } @@ -7412,18 +7462,14 @@ static void *session_do(void *data) 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); @@ -7511,33 +7557,52 @@ static int purge_sessions(int n_max) 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) @@ -7556,14 +7621,10 @@ 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, @@ -7572,86 +7633,89 @@ static int __attribute__((format(printf, 9, 0))) __manager_event_sessions_va( 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); @@ -7662,11 +7726,13 @@ static int __attribute__((format(printf, 9, 0))) __manager_event_sessions_va( 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; } @@ -8423,13 +8489,12 @@ static int generic_http_callback(struct ast_tcptls_session_instance *ser, /* 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 @@ -8439,8 +8504,6 @@ static int generic_http_callback(struct ast_tcptls_session_instance *ser, */ while ((session->managerid = ast_random() ^ (unsigned long) session) == 0) { } - session->last_ev = grab_last(); - AST_LIST_HEAD_INIT_NOLOCK(&session->datastores); } ao2_unlock(session); @@ -8563,11 +8626,11 @@ static int generic_http_callback(struct ast_tcptls_session_instance *ser, 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 { @@ -8739,7 +8802,7 @@ static int auth_http_callback(struct ast_tcptls_session_instance *ser, * 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; @@ -8748,8 +8811,6 @@ static int auth_http_callback(struct ast_tcptls_session_instance *ser, 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; @@ -9134,7 +9195,6 @@ static void purge_old_stuff(void *data) } else { ser->poll_timeout = 5000; } - purge_events(); } static struct ast_tls_config ami_tls_cfg; @@ -9418,7 +9478,6 @@ static struct ast_cli_entry cli_manager[] = { 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"), @@ -9781,12 +9840,6 @@ static int __init_manager(int reload, int by_external_config) __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) {