diff --git a/ChangeLog b/ChangeLog
index bb320ab..d476c2d 100644
--- a/ChangeLog
+++ b/ChangeLog
@@ -2,9 +2,13 @@
- Fix potential infinite loop in resolver (#200)
- Add code coverage support
- Add support for password-protected TLS key & PKCS#12/PFX files
+ - Stream-Management support (XEP-0198)
- New API:
- xmpp_conn_send_queue_len()
- xmpp_conn_send_queue_drop_element()
+ - xmpp_conn_get_sm_state()
+ - xmpp_conn_set_sm_state()
+ - xmpp_free_sm_state()
- xmpp_conn_get_keyfile()
- xmpp_conn_set_password_callback()
- xmpp_conn_set_password_retries()
diff --git a/src/auth.c b/src/auth.c
index fcc57ad..dd1e541 100644
--- a/src/auth.c
+++ b/src/auth.c
@@ -879,21 +879,79 @@ static void _handle_open_sasl(xmpp_conn_t *conn)
NULL);
}
+static int _do_bind(xmpp_conn_t *conn, xmpp_stanza_t *bind)
+{
+ xmpp_stanza_t *iq, *res, *text;
+ char *resource;
+
+ /* setup response handlers */
+ handler_add_id(conn, _handle_bind, "_xmpp_bind1", NULL);
+ handler_add_timed(conn, _handle_missing_bind, BIND_TIMEOUT, NULL);
+
+ /* send bind request */
+ iq = xmpp_iq_new(conn->ctx, "set", "_xmpp_bind1");
+ if (!iq) {
+ disconnect_mem_error(conn);
+ return 0;
+ }
+
+ /* request a specific resource if we have one */
+ resource = xmpp_jid_resource(conn->ctx, conn->jid);
+ if ((resource != NULL) && (strlen(resource) == 0)) {
+ /* jabberd2 doesn't handle an empty resource */
+ strophe_free(conn->ctx, resource);
+ resource = NULL;
+ }
+
+ /* if we have a resource to request, do it. otherwise the
+ server will assign us one */
+ if (resource) {
+ res = xmpp_stanza_new(conn->ctx);
+ if (!res) {
+ xmpp_stanza_release(bind);
+ xmpp_stanza_release(iq);
+ disconnect_mem_error(conn);
+ return 0;
+ }
+ xmpp_stanza_set_name(res, "resource");
+ text = xmpp_stanza_new(conn->ctx);
+ if (!text) {
+ xmpp_stanza_release(res);
+ xmpp_stanza_release(bind);
+ xmpp_stanza_release(iq);
+ disconnect_mem_error(conn);
+ return 0;
+ }
+ xmpp_stanza_set_text(text, resource);
+ xmpp_stanza_add_child(res, text);
+ xmpp_stanza_release(text);
+ xmpp_stanza_add_child(bind, res);
+ xmpp_stanza_release(res);
+ strophe_free(conn->ctx, resource);
+ }
+
+ xmpp_stanza_add_child(iq, bind);
+ xmpp_stanza_release(bind);
+
+ /* send bind request */
+ send_stanza(conn, iq, XMPP_QUEUE_STROPHE);
+ xmpp_stanza_release(iq);
+ return 0;
+}
+
static int
_handle_features_sasl(xmpp_conn_t *conn, xmpp_stanza_t *stanza, void *userdata)
{
- xmpp_stanza_t *bind, *session, *iq, *res, *text, *opt;
+ xmpp_stanza_t *bind, *session, *opt;
+ xmpp_stanza_t *resume;
const char *ns;
- char *resource;
+ char h[11];
UNUSED(userdata);
/* remove missing features handler */
xmpp_timed_handler_delete(conn, _handle_missing_features_sasl);
- /* we are expecting and since this is a
- XMPP style connection */
-
/* check whether resource binding is required */
bind = xmpp_stanza_get_child_by_name(stanza, "bind");
if (bind) {
@@ -911,77 +969,46 @@ _handle_features_sasl(xmpp_conn_t *conn, xmpp_stanza_t *stanza, void *userdata)
ns != NULL && strcmp(ns, XMPP_NS_SESSION) == 0;
}
- opt = xmpp_stanza_get_child_by_ns(stanza, XMPP_NS_SM);
- if (opt && strcmp(xmpp_stanza_get_name(opt), "sm") == 0) {
+ if (xmpp_stanza_get_child_by_name_and_ns(stanza, "sm", XMPP_NS_SM)) {
/* stream management supported */
- conn->sm_support = 1;
+ conn->sm_state->sm_support = 1;
}
+ bind = xmpp_stanza_copy(bind);
+ if (!bind) {
+ disconnect_mem_error(conn);
+ return 0;
+ }
+
+ /* we are expecting either and since this is a
+ XMPP style connection or we the previous session */
+
+ /* check whether we can the previous session */
+ if (!conn->sm_disable && conn->sm_state->can_resume &&
+ conn->sm_state->previd && conn->sm_state->bound_jid) {
+ resume = xmpp_stanza_new(conn->ctx);
+ if (!resume) {
+ disconnect_mem_error(conn);
+ return 0;
+ }
+ conn->sm_state->bind = bind;
+ conn->sm_state->resume = 1;
+ xmpp_stanza_set_name(resume, "resume");
+ xmpp_stanza_set_ns(resume, XMPP_NS_SM);
+ xmpp_stanza_set_attribute(resume, "previd", conn->sm_state->previd);
+ strophe_snprintf(h, sizeof(h), "%u", conn->sm_state->sm_handled_nr);
+ xmpp_stanza_set_attribute(resume, "h", h);
+ send_stanza(conn, resume, XMPP_QUEUE_SM_STROPHE);
+ xmpp_stanza_release(resume);
+ handler_add(conn, _handle_sm, XMPP_NS_SM, NULL, NULL, NULL);
+ }
/* if bind is required, go ahead and start it */
- if (conn->bind_required) {
+ else if (conn->bind_required) {
/* bind resource */
-
- /* setup response handlers */
- handler_add_id(conn, _handle_bind, "_xmpp_bind1", NULL);
- handler_add_timed(conn, _handle_missing_bind, BIND_TIMEOUT, NULL);
-
- /* send bind request */
- iq = xmpp_iq_new(conn->ctx, "set", "_xmpp_bind1");
- if (!iq) {
- disconnect_mem_error(conn);
- return 0;
- }
-
- bind = xmpp_stanza_copy(bind);
- if (!bind) {
- xmpp_stanza_release(iq);
- disconnect_mem_error(conn);
- return 0;
- }
-
- /* request a specific resource if we have one */
- resource = xmpp_jid_resource(conn->ctx, conn->jid);
- if ((resource != NULL) && (strlen(resource) == 0)) {
- /* jabberd2 doesn't handle an empty resource */
- strophe_free(conn->ctx, resource);
- resource = NULL;
- }
-
- /* if we have a resource to request, do it. otherwise the
- server will assign us one */
- if (resource) {
- res = xmpp_stanza_new(conn->ctx);
- if (!res) {
- xmpp_stanza_release(bind);
- xmpp_stanza_release(iq);
- disconnect_mem_error(conn);
- return 0;
- }
- xmpp_stanza_set_name(res, "resource");
- text = xmpp_stanza_new(conn->ctx);
- if (!text) {
- xmpp_stanza_release(res);
- xmpp_stanza_release(bind);
- xmpp_stanza_release(iq);
- disconnect_mem_error(conn);
- return 0;
- }
- xmpp_stanza_set_text(text, resource);
- xmpp_stanza_add_child(res, text);
- xmpp_stanza_release(text);
- xmpp_stanza_add_child(bind, res);
- xmpp_stanza_release(res);
- strophe_free(conn->ctx, resource);
- }
-
- xmpp_stanza_add_child(iq, bind);
- xmpp_stanza_release(bind);
-
- /* send bind request */
- send_stanza(conn, iq, XMPP_QUEUE_STROPHE);
- xmpp_stanza_release(iq);
+ _do_bind(conn, bind);
} else {
/* can't bind, disconnect */
+ xmpp_stanza_release(bind);
strophe_error(conn->ctx, "xmpp",
"Stream features does not allow "
"resource bind.");
@@ -1061,7 +1088,7 @@ _handle_bind(xmpp_conn_t *conn, xmpp_stanza_t *stanza, void *userdata)
xmpp_stanza_release(iq);
}
- if (conn->sm_support) {
+ if (conn->sm_state->sm_support && !conn->sm_disable) {
enable = xmpp_stanza_new(conn->ctx);
if (!enable) {
disconnect_mem_error(conn);
@@ -1069,10 +1096,11 @@ _handle_bind(xmpp_conn_t *conn, xmpp_stanza_t *stanza, void *userdata)
}
xmpp_stanza_set_name(enable, "enable");
xmpp_stanza_set_ns(enable, XMPP_NS_SM);
+ if (!conn->sm_state->dont_request_resume)
+ xmpp_stanza_set_attribute(enable, "resume", "true");
handler_add(conn, _handle_sm, XMPP_NS_SM, NULL, NULL, NULL);
send_stanza(conn, enable, XMPP_QUEUE_SM_STROPHE);
xmpp_stanza_release(enable);
- conn->sm_sent_nr = 0;
}
if (!conn->session_required) {
@@ -1148,15 +1176,129 @@ static int _handle_sm(xmpp_conn_t *const conn,
xmpp_stanza_t *const stanza,
void *const userdata)
{
- const char *name;
+ xmpp_stanza_t *failed_cause;
+ const char *name, *id, *previd, *resume, *h, *cause;
+ xmpp_send_queue_t *e;
+ unsigned long ul_h = 0;
+
+ UNUSED(userdata);
name = xmpp_stanza_get_name(stanza);
- if (name && strcmp(name, "enabled") == 0) {
- conn->sm_enabled = 1;
- conn->sm_handled_nr = 0;
+ if (!name)
+ goto LBL_ERR;
+
+ if (strcmp(name, "enabled") == 0) {
+ conn->sm_state->sm_enabled = 1;
+ conn->sm_state->sm_handled_nr = 0;
+ resume = xmpp_stanza_get_attribute(stanza, "resume");
+ if (resume && (strcasecmp(resume, "true") || strcmp(resume, "1"))) {
+ id = xmpp_stanza_get_attribute(stanza, "id");
+ if (!id) {
+ strophe_error(conn->ctx, "xmpp",
+ "SM error: server said it can resume, but "
+ "didn't provide an ID.");
+ name = NULL;
+ goto LBL_ERR;
+ }
+ conn->sm_state->can_resume = 1;
+ conn->sm_state->id = strophe_strdup(conn->ctx, id);
+ }
+ } else if (strcmp(name, "resumed") == 0) {
+ previd = xmpp_stanza_get_attribute(stanza, "previd");
+ if (!previd || strcmp(previd, conn->sm_state->previd)) {
+ strophe_error(conn->ctx, "xmpp",
+ "SM error: previd didn't match, ours is \"%s\".",
+ conn->sm_state->previd);
+ name = NULL;
+ goto LBL_ERR;
+ }
+ h = xmpp_stanza_get_attribute(stanza, "h");
+ if (!h || string_to_ul(h, &ul_h)) {
+ strophe_error(conn->ctx, "xmpp",
+ "SM error: failed parsing 'h', it got converted "
+ "to %llu.",
+ ul_h);
+ name = NULL;
+ goto LBL_ERR;
+ }
+ conn->sm_state->sm_enabled = 1;
+ conn->sm_state->id = conn->sm_state->previd;
+ conn->sm_state->previd = NULL;
+ conn->bound_jid = conn->sm_state->bound_jid;
+ conn->sm_state->bound_jid = NULL;
+ if (conn->sm_state->sm_queue.head)
+ conn->sm_state->sm_sent_nr = conn->sm_state->sm_queue.head->sm_h;
+ else
+ conn->sm_state->sm_sent_nr = ul_h;
+ while ((e = pop_queue_front(&conn->sm_state->sm_queue))) {
+ if (e->sm_h >= ul_h) {
+ /* Re-send what was already sent out and is still in the
+ * SM queue (i.e. it hasn't been ACK'ed by the server)
+ */
+ send_raw(conn, e->data, e->len, e->owner, NULL);
+ }
+ strophe_free(conn->ctx, queue_element_free(conn->ctx, e));
+ }
+ conn->authenticated = 1;
+
+ /* call connection handler */
+ conn->conn_handler(conn, XMPP_CONN_CONNECT, 0, NULL, conn->userdata);
+ } else if (strcmp(name, "failed") == 0) {
+ name = NULL;
+
+ failed_cause =
+ xmpp_stanza_get_child_by_ns(stanza, XMPP_NS_STANZAS_IETF);
+ if (!failed_cause)
+ goto LBL_ERR;
+
+ cause = xmpp_stanza_get_name(failed_cause);
+ if (!cause)
+ goto LBL_ERR;
+
+ if (!strcmp(cause, "item-not-found") ||
+ !strcmp(cause, "feature-not-implemented")) {
+ if (conn->sm_state->resume) {
+ conn->sm_state->resume = 0;
+ conn->sm_state->can_resume = 0;
+ /* remember that the server reports having support
+ * for resumption, but actually it doesn't ...
+ */
+ conn->sm_state->dont_request_resume =
+ !strcmp(cause, "feature-not-implemented");
+ strophe_free(conn->ctx, conn->sm_state->previd);
+ conn->sm_state->previd = NULL;
+ strophe_free(conn->ctx, conn->sm_state->bound_jid);
+ conn->sm_state->bound_jid = NULL;
+ _do_bind(conn, conn->sm_state->bind);
+ conn->sm_state->bind = NULL;
+ }
+ }
+ conn->sm_state->sm_handled_nr = 0;
} else {
- conn->sm_enabled = 0;
- strophe_warn(conn->ctx, "auth", "Stream management failed.");
+ /* unknown stanza received */
+ name = NULL;
+ }
+
+LBL_ERR:
+ if (!name) {
+ char *err = "Couldn't convert stanza to text!";
+ char *buf;
+ size_t buflen;
+ switch (xmpp_stanza_to_text(stanza, &buf, &buflen)) {
+ case XMPP_EOK:
+ break;
+ case XMPP_EMEM:
+ disconnect_mem_error(conn);
+ return 0;
+ default:
+ buf = err;
+ break;
+ }
+ strophe_warn(conn->ctx, "xmpp", "SM error: Stanza received was: %s",
+ buf);
+ if (buf != err)
+ strophe_free(conn->ctx, buf);
+ conn->sm_state->sm_enabled = 0;
}
return 0;
}
diff --git a/src/common.h b/src/common.h
index 5723484..8ff7438 100644
--- a/src/common.h
+++ b/src/common.h
@@ -140,9 +140,12 @@ struct _xmpp_send_queue_t {
char *data;
size_t len;
size_t written;
+ int wip;
xmpp_send_queue_owner_t owner;
+ void *userdata;
+ uint32_t sm_h;
- xmpp_send_queue_t *next;
+ xmpp_send_queue_t *prev, *next;
};
#define UNUSED(x) ((void)(x))
@@ -168,6 +171,22 @@ enum {
typedef void (*xmpp_open_handler)(xmpp_conn_t *conn);
+typedef struct {
+ xmpp_send_queue_t *head, *tail;
+} xmpp_queue_t;
+
+struct _xmpp_sm_t {
+ xmpp_ctx_t *ctx;
+ int sm_support;
+ int sm_enabled;
+ int can_resume, resume, dont_request_resume;
+ uint32_t sm_handled_nr;
+ uint32_t sm_sent_nr;
+ xmpp_queue_t sm_queue;
+ char *id, *previd, *bound_jid;
+ xmpp_stanza_t *bind;
+};
+
struct _xmpp_conn_t {
unsigned int ref;
xmpp_ctx_t *ctx;
@@ -212,10 +231,8 @@ struct _xmpp_conn_t {
/* if server returns or we must do them */
int bind_required;
int session_required;
- int sm_support;
- int sm_enabled;
- uint32_t sm_handled_nr;
- uint32_t sm_sent_nr;
+ int sm_disable;
+ xmpp_sm_state_t *sm_state;
char *lang;
char *domain;
@@ -317,11 +334,17 @@ void auth_handle_component_open(xmpp_conn_t *conn);
void auth_handle_open_raw(xmpp_conn_t *conn);
void auth_handle_open_stub(xmpp_conn_t *conn);
+/* queue functions */
+void add_queue_back(xmpp_queue_t *queue, xmpp_send_queue_t *item);
+xmpp_send_queue_t *pop_queue_front(xmpp_queue_t *queue);
+char *queue_element_free(xmpp_ctx_t *ctx, xmpp_send_queue_t *e);
+
/* send functions */
void send_raw(xmpp_conn_t *conn,
const char *data,
size_t len,
- xmpp_send_queue_owner_t owner);
+ xmpp_send_queue_owner_t owner,
+ void *userdata);
/* this is a bit special as it will always mark the sent string as
* owned by libstrophe
*/
diff --git a/src/conn.c b/src/conn.c
index c00d9a4..645884c 100644
--- a/src/conn.c
+++ b/src/conn.c
@@ -22,7 +22,6 @@
#include
#include
#include
-#include
#include
#include "strophe.h"
@@ -100,7 +99,6 @@ static void _conn_sm_handle_stanza(xmpp_conn_t *const conn,
xmpp_stanza_t *stanza);
static unsigned short _conn_default_port(xmpp_conn_t *conn,
xmpp_conn_type_t type);
-static char *_queue_element_free(xmpp_ctx_t *ctx, xmpp_send_queue_t *e);
static void _conn_reset(xmpp_conn_t *conn);
static int _conn_connect(xmpp_conn_t *conn,
const char *domain,
@@ -116,7 +114,8 @@ static void _send_valist(xmpp_conn_t *conn,
static int _send_raw(xmpp_conn_t *conn,
char *data,
size_t len,
- xmpp_send_queue_owner_t owner);
+ xmpp_send_queue_owner_t owner,
+ void *userdata);
void xmpp_send_error(xmpp_conn_t *conn, xmpp_error_type_t type, char *text)
{
@@ -149,6 +148,7 @@ xmpp_conn_t *xmpp_conn_new(xmpp_ctx_t *ctx)
conn->type = XMPP_UNKNOWN;
conn->state = XMPP_STATE_DISCONNECTED;
+
conn->sock = -1;
conn->ka_timeout = KEEPALIVE_TIMEOUT;
conn->ka_interval = KEEPALIVE_INTERVAL;
@@ -202,10 +202,7 @@ xmpp_conn_t *xmpp_conn_new(xmpp_ctx_t *ctx)
conn->bind_required = 0;
conn->session_required = 0;
- conn->sm_support = 0;
- conn->sm_enabled = 0;
- conn->sm_handled_nr = 0;
- conn->sm_sent_nr = 0;
+ conn->sm_state = NULL;
conn->parser =
parser_new(conn->ctx, _handle_stream_start, _handle_stream_end,
@@ -403,6 +400,8 @@ int xmpp_conn_release(xmpp_conn_t *conn)
strophe_free(ctx, conn->tls_cafile);
if (conn->tls_capath)
strophe_free(ctx, conn->tls_capath);
+ if (conn->sm_state)
+ xmpp_free_sm_state(conn->sm_state);
tls_clear_password_cache(conn);
strophe_free(ctx, conn);
released = 1;
@@ -725,6 +724,14 @@ int xmpp_connect_client(xmpp_conn_t *conn,
if (!domain)
return XMPP_EMEM;
+ if (!conn->sm_state) {
+ conn->sm_state = strophe_alloc(conn->ctx, sizeof(*conn->sm_state));
+ if (!conn->sm_state)
+ return XMPP_EMEM;
+ memset(conn->sm_state, 0, sizeof(*conn->sm_state));
+ conn->sm_state->ctx = conn->ctx;
+ }
+
if (altdomain != NULL) {
strophe_debug(conn->ctx, "xmpp", "Connecting via altdomain.");
host = altdomain;
@@ -1041,7 +1048,7 @@ void xmpp_send_raw_string(xmpp_conn_t *conn, const char *fmt, ...)
*/
void xmpp_send_raw(xmpp_conn_t *conn, const char *data, size_t len)
{
- send_raw(conn, data, len, XMPP_QUEUE_USER);
+ send_raw(conn, data, len, XMPP_QUEUE_USER, NULL);
}
/** Send an XML stanza to the XMPP server.
@@ -1189,6 +1196,7 @@ int xmpp_conn_set_flags(xmpp_conn_t *conn, long flags)
conn->tls_legacy_ssl = (flags & XMPP_CONN_FLAG_LEGACY_SSL) ? 1 : 0;
conn->tls_trust = (flags & XMPP_CONN_FLAG_TRUST_TLS) ? 1 : 0;
conn->auth_legacy_enabled = (flags & XMPP_CONN_FLAG_LEGACY_AUTH) ? 1 : 0;
+ conn->sm_disable = (flags & XMPP_CONN_FLAG_DISABLE_SM) ? 1 : 0;
return 0;
}
@@ -1253,6 +1261,124 @@ int xmpp_conn_is_disconnected(xmpp_conn_t *conn)
return conn->state == XMPP_STATE_DISCONNECTED;
}
+/**
+ * This returns the Stream Management state of a connection object after
+ * it has been disconnected.
+ * One can then initialise a fresh connection object and set this Stream
+ * Management state by calling \ref xmpp_conn_set_sm_state
+ *
+ * In case one wants to dispose of the state w/o setting it into a fresh
+ * connection object, one can call \ref xmpp_free_sm_state
+ *
+ * After calling this function to retrieve the state, only call one of the
+ * other two.
+ *
+ * @param conn a Strophe connection object
+ * @return The Stream Management state of the connection or NULL on error
+ *
+ * @ingroup Connections
+ */
+xmpp_sm_state_t *xmpp_conn_get_sm_state(xmpp_conn_t *conn)
+{
+ xmpp_sm_state_t *ret;
+
+ /* We can only return the SM state when we're disconnected */
+ if (conn->state != XMPP_STATE_DISCONNECTED)
+ return NULL;
+
+ ret = conn->sm_state;
+ conn->sm_state = NULL;
+ if (ret->previd) {
+ strophe_free(conn->ctx, ret->previd);
+ ret->previd = NULL;
+ }
+
+ if (ret->can_resume) {
+ ret->previd = ret->id;
+ ret->id = NULL;
+
+ ret->bound_jid = conn->bound_jid;
+ conn->bound_jid = NULL;
+ } else if (ret->id) {
+ strophe_free(conn->ctx, ret->id);
+ ret->id = NULL;
+ }
+
+ ret->sm_enabled = ret->sm_support = ret->resume = 0;
+
+ if (ret->bind) {
+ xmpp_stanza_release(ret->bind);
+ ret->bind = NULL;
+ }
+
+ return ret;
+}
+
+/**
+ * @param conn a Strophe connection object
+ * @param sm_state A Stream Management state returned from a call to
+ * `xmpp_conn_get_sm_state()`
+ *
+ * @return XMPP_EOK (0) on success or a number less than 0 on failure
+ *
+ * @ingroup Connections
+ */
+int xmpp_conn_set_sm_state(xmpp_conn_t *conn, xmpp_sm_state_t *sm_state)
+{
+ /* We can only set the SM state when we're disconnected */
+ if (conn->state != XMPP_STATE_DISCONNECTED)
+ return XMPP_EINVOP;
+
+ if (conn->sm_state) {
+ strophe_error(conn->ctx, "conn", "SM state is already set!");
+ return XMPP_EINVOP;
+ }
+
+ if (conn->ctx != sm_state->ctx) {
+ strophe_error(
+ conn->ctx, "conn",
+ "SM state has to be assigned to connection that stems from "
+ "the same context!");
+ return XMPP_EINVOP;
+ }
+
+ conn->sm_state = sm_state;
+ return XMPP_EOK;
+}
+
+/** c.f. \ref xmpp_conn_get_sm_state for usage documentation
+ *
+ * @param sm_state A Stream Management state returned from a call to
+ * `xmpp_conn_get_sm_state()`
+ *
+ * @ingroup Connections
+ */
+void xmpp_free_sm_state(xmpp_sm_state_t *sm_state)
+{
+ xmpp_ctx_t *ctx;
+ xmpp_send_queue_t *smq;
+
+ if (!sm_state || !sm_state->ctx)
+ return;
+
+ ctx = sm_state->ctx;
+
+ while ((smq = pop_queue_front(&sm_state->sm_queue))) {
+ strophe_free(ctx, queue_element_free(ctx, smq));
+ }
+
+ if (sm_state->bind)
+ xmpp_stanza_release(sm_state->bind);
+ if (sm_state->id)
+ strophe_free(ctx, sm_state->id);
+ if (sm_state->previd)
+ strophe_free(ctx, sm_state->previd);
+ if (sm_state->bound_jid)
+ strophe_free(ctx, sm_state->bound_jid);
+ memset(sm_state, 0, sizeof(*sm_state));
+ strophe_free(ctx, sm_state);
+}
+
/**
* @return The number of entries in the send queue
*
@@ -1260,82 +1386,99 @@ int xmpp_conn_is_disconnected(xmpp_conn_t *conn)
*/
int xmpp_conn_send_queue_len(const xmpp_conn_t *conn)
{
- if (conn->send_queue_head && conn->send_queue_head->written &&
+ if (conn->send_queue_head && conn->send_queue_head->wip &&
conn->send_queue_head->owner == XMPP_QUEUE_USER)
return conn->send_queue_user_len - 1;
else
return conn->send_queue_user_len;
}
+static char *_drop_send_queue_element(xmpp_conn_t *conn, xmpp_send_queue_t *e)
+{
+ if (e == conn->send_queue_head)
+ conn->send_queue_head = e->next;
+ if (e == conn->send_queue_tail)
+ conn->send_queue_tail = e->prev;
+ if (!conn->send_queue_head)
+ conn->send_queue_tail = NULL;
+ if (e->prev)
+ e->prev->next = e->next;
+ if (e->next)
+ e->next->prev = e->prev;
+ conn->send_queue_len--;
+ if (e->owner == XMPP_QUEUE_USER)
+ conn->send_queue_user_len--;
+ return queue_element_free(conn->ctx, e);
+}
+
/** Drop an element of the send queue.
* This can be used to manage the send queue in case a server
* isn't fast enough in processing the elements you're trying
- * to send.
+ * to send or your outgoing bandwidth isn't fast enough to transfer
+ * everything you want to send out.
*
* @param conn a Strophe connection object
* @param which the element that shall be removed
*
+ * @return The rendered stanza. The pointer returned has to be free'd by the
+ * caller of this function.
+ *
* @ingroup Connections
*/
char *xmpp_conn_send_queue_drop_element(xmpp_conn_t *conn,
xmpp_queue_element_t which)
{
- xmpp_send_queue_t *t, *p;
- xmpp_send_queue_owner_t owner;
- char *ret;
+ xmpp_send_queue_t *t;
+
+ /* Fast return paths */
/* empty queue */
if (!conn->send_queue_head)
return NULL;
/* one element in queue */
if (conn->send_queue_head == conn->send_queue_tail) {
- if (conn->send_queue_head->written)
- return NULL;
-
- t = conn->send_queue_head;
- conn->send_queue_head = conn->send_queue_tail = NULL;
- owner = t->owner;
- ret = _queue_element_free(conn->ctx, t);
- conn->send_queue_len--;
- if (owner == XMPP_QUEUE_USER)
- conn->send_queue_user_len--;
- return ret;
- }
- switch (which) {
- case XMPP_QUEUE_OLDEST:
/* head is already sent out partially */
- if (conn->send_queue_head->written) {
- t = conn->send_queue_head->next;
- /* there are no more elements in the queue */
- if (!t)
- return NULL;
- conn->send_queue_head->next = t->next;
- } else {
- t = conn->send_queue_head;
- conn->send_queue_head = t->next;
- }
- owner = t->owner;
- ret = _queue_element_free(conn->ctx, t);
- conn->send_queue_len--;
- if (owner == XMPP_QUEUE_USER)
- conn->send_queue_user_len--;
- return ret;
- case XMPP_QUEUE_YOUNGEST:
+ if (conn->send_queue_head->wip)
+ return NULL;
+ /* the element is no USER element */
+ if (conn->send_queue_head->owner != XMPP_QUEUE_USER)
+ return NULL;
+ }
+
+ /* Regular flow */
+ if (which == XMPP_QUEUE_OLDEST) {
t = conn->send_queue_head;
- do {
- p = t;
- t = t->next;
- } while (t != conn->send_queue_tail);
- conn->send_queue_tail = p;
- owner = t->owner;
- ret = _queue_element_free(conn->ctx, t);
- conn->send_queue_len--;
- if (owner == XMPP_QUEUE_USER)
- conn->send_queue_user_len--;
- return ret;
- default:
+ } else if (which == XMPP_QUEUE_YOUNGEST) {
+ t = conn->send_queue_tail;
+ /* search backwards to find last USER element */
+ while (t && t->owner != XMPP_QUEUE_USER)
+ t = t->prev;
+ } else {
strophe_error(conn->ctx, "conn", "Unknown queue element %d", which);
return NULL;
}
+ /* there was no USER element in the queue */
+ if (!t)
+ return NULL;
+
+ /* head is already sent out partially */
+ if (t == conn->send_queue_head && t->wip)
+ t = t->next;
+
+ /* search forward to find the first USER element */
+ while (t && t->owner != XMPP_QUEUE_USER)
+ t = t->next;
+
+ /* there was no USER element in the queue we could drop */
+ if (!t)
+ return NULL;
+
+ /* In case there exists a SM stanza that is linked to the
+ * one we're currently dropping, also delete that one.
+ */
+ if (t->next && t->next->userdata == t)
+ strophe_free(conn->ctx, _drop_send_queue_element(conn, t->next));
+ /* Finally drop the element */
+ return _drop_send_queue_element(conn, t);
}
/* timed handler for cleanup if normal disconnect procedure takes too long */
@@ -1527,6 +1670,8 @@ static void _handle_stream_end(char *name, void *userdata)
/* stream is over */
strophe_debug(conn->ctx, "xmpp", "RECV: ");
+ /* the session has been terminated properly, i.e. it can't be resumed */
+ conn->sm_state->can_resume = 0;
conn_disconnect_clean(conn);
}
@@ -1542,7 +1687,7 @@ static void _handle_stream_stanza(xmpp_stanza_t *stanza, void *userdata)
}
handler_fire_stanza(conn, stanza);
- if (conn->sm_enabled)
+ if (conn->sm_state->sm_enabled)
_conn_sm_handle_stanza(conn, stanza);
}
@@ -1551,25 +1696,54 @@ static void _conn_sm_handle_stanza(xmpp_conn_t *const conn,
xmpp_stanza_t *stanza)
{
xmpp_stanza_t *a;
- const char *name;
- const char *ns;
+ xmpp_send_queue_t *e;
+ char *c;
+ const char *name, *ns, *attr_h;
char h[11];
+ unsigned long ul_h;
ns = xmpp_stanza_get_ns(stanza);
if (ns && strcmp(ns, XMPP_NS_SM) != 0)
- ++conn->sm_handled_nr;
+ ++conn->sm_state->sm_handled_nr;
else {
name = xmpp_stanza_get_name(stanza);
- if (name && strcmp(name, "r") == 0) {
+ if (!name)
+ return;
+ if (strcmp(name, "r") == 0) {
a = xmpp_stanza_new(conn->ctx);
- if (!a)
- return; /* XXX */
+ if (!a) {
+ strophe_debug(conn->ctx, "conn", "Couldn't create stanza.");
+ return;
+ }
xmpp_stanza_set_name(a, "a");
xmpp_stanza_set_ns(a, XMPP_NS_SM);
- strophe_snprintf(h, sizeof(h), "%u", conn->sm_handled_nr);
+ strophe_snprintf(h, sizeof(h), "%u", conn->sm_state->sm_handled_nr);
xmpp_stanza_set_attribute(a, "h", h);
send_stanza(conn, a, XMPP_QUEUE_SM_STROPHE);
xmpp_stanza_release(a);
+ } else if (strcmp(name, "a") == 0) {
+ attr_h = xmpp_stanza_get_attribute(stanza, "h");
+ if (!attr_h) {
+ strophe_debug(conn->ctx, "conn", "Didn't find 'h' attribute.");
+ return;
+ }
+ if (string_to_ul(attr_h, &ul_h)) {
+ strophe_error(
+ conn->ctx, "conn",
+ "Error on strtoul() of '%s', returned value is %llu.",
+ attr_h, ul_h);
+ /* We continue here and drop the complete SM queue instead of
+ * returning and letting the queue fill up.
+ */
+ ul_h = ULONG_MAX;
+ }
+ while (conn->sm_state->sm_queue.head &&
+ conn->sm_state->sm_queue.head->sm_h < ul_h) {
+ e = pop_queue_front(&conn->sm_state->sm_queue);
+ strophe_debug_verbose(2, conn->ctx, "conn", "SM_Q_DROP: %p", e);
+ c = queue_element_free(conn->ctx, e);
+ strophe_free(conn->ctx, c);
+ }
}
}
}
@@ -1588,7 +1762,7 @@ static unsigned short _conn_default_port(xmpp_conn_t *conn,
};
}
-static char *_queue_element_free(xmpp_ctx_t *ctx, xmpp_send_queue_t *e)
+char *queue_element_free(xmpp_ctx_t *ctx, xmpp_send_queue_t *e)
{
char *ret = e->data;
strophe_debug_verbose(2, ctx, "conn", "Q_FREE: %p", e);
@@ -1613,7 +1787,7 @@ static void _conn_reset(xmpp_conn_t *conn)
while (sq) {
tsq = sq;
sq = sq->next;
- strophe_free(ctx, _queue_element_free(ctx, tsq));
+ strophe_free(ctx, queue_element_free(ctx, tsq));
}
conn->send_queue_head = NULL;
conn->send_queue_tail = NULL;
@@ -1643,8 +1817,7 @@ static void _conn_reset(xmpp_conn_t *conn)
conn->error = 0;
conn->tls_support = 0;
- conn->sm_support = 0;
- conn->sm_enabled = 0;
+
conn->bind_required = 0;
conn->session_required = 0;
@@ -1705,7 +1878,8 @@ static int _conn_connect(xmpp_conn_t *conn,
void send_raw(xmpp_conn_t *conn,
const char *data,
size_t len,
- xmpp_send_queue_owner_t owner)
+ xmpp_send_queue_owner_t owner,
+ void *userdata)
{
char *d;
@@ -1718,7 +1892,7 @@ void send_raw(xmpp_conn_t *conn,
return;
}
- _send_raw(conn, d, len, owner);
+ _send_raw(conn, d, len, owner, userdata);
}
static void _send_valist(xmpp_conn_t *conn,
@@ -1753,10 +1927,10 @@ static void _send_valist(xmpp_conn_t *conn,
va_end(apdup);
/* len - 1 so we don't send trailing \0 */
- _send_raw(conn, bigbuf, len - 1, owner);
+ _send_raw(conn, bigbuf, len - 1, owner, NULL);
} else {
/* go through send_raw() which does the strdup() for us */
- send_raw(conn, buf, len, owner);
+ send_raw(conn, buf, len, owner, NULL);
}
}
@@ -1787,13 +1961,43 @@ void send_stanza(xmpp_conn_t *conn,
return;
}
- _send_raw(conn, buf, len, owner);
+ _send_raw(conn, buf, len, owner, NULL);
+}
+
+void add_queue_back(xmpp_queue_t *queue, xmpp_send_queue_t *item)
+{
+ item->next = NULL;
+ if (!queue->tail) {
+ item->prev = NULL;
+ queue->head = item;
+ queue->tail = item;
+ } else {
+ item->prev = queue->tail;
+ queue->tail->next = item;
+ queue->tail = item;
+ }
+}
+
+xmpp_send_queue_t *pop_queue_front(xmpp_queue_t *queue)
+{
+ xmpp_send_queue_t *ret = queue->head;
+ if (queue->head) {
+ queue->head = queue->head->next;
+ if (!queue->head) {
+ queue->tail = NULL;
+ } else {
+ queue->head->prev = NULL;
+ }
+ ret->prev = ret->next = NULL;
+ }
+ return ret;
}
static int _send_raw(xmpp_conn_t *conn,
char *data,
size_t len,
- xmpp_send_queue_owner_t owner)
+ xmpp_send_queue_owner_t owner,
+ void *userdata)
{
xmpp_send_queue_t *item;
const char *req_ack = "";
@@ -1809,7 +2013,10 @@ static int _send_raw(xmpp_conn_t *conn,
item->data = data;
item->len = len;
item->next = NULL;
+ item->prev = conn->send_queue_tail;
item->written = 0;
+ item->wip = 0;
+ item->userdata = userdata;
item->owner = owner;
if (!conn->send_queue_tail) {
@@ -1826,5 +2033,8 @@ static int _send_raw(xmpp_conn_t *conn,
conn->send_queue_user_len++;
strophe_debug_verbose(3, conn->ctx, "conn", "QUEUED: %s", data);
strophe_debug_verbose(1, conn->ctx, "conn", "Q_ADD: %p", item);
+ if (!(owner & XMPP_QUEUE_SM) && conn->sm_state->sm_enabled) {
+ send_raw(conn, req_ack, strlen(req_ack), XMPP_QUEUE_SM_STROPHE, item);
+ }
return XMPP_EOK;
}
diff --git a/src/event.c b/src/event.c
index 3bf4a73..953d7f3 100644
--- a/src/event.c
+++ b/src/event.c
@@ -131,20 +131,32 @@ void xmpp_run_once(xmpp_ctx_t *ctx, unsigned long timeout)
}
if (ret > 0 && ret < towrite)
sq->written += ret; /* not all data could be sent now */
+ sq->wip = 1;
if (ret != towrite)
break; /* partial write or an error */
/* all data for this queue item written, delete and move on */
strophe_debug(conn->ctx, "conn", "SENT: %s", sq->data);
- strophe_debug_verbose(1, ctx, "xmpp",
- "Finished writing queue element: %p.", sq);
- strophe_free(ctx, sq->data);
+ strophe_debug_verbose(1, ctx, "xmpp", "Q_SENT: %p", sq);
tsq = sq;
sq = sq->next;
conn->send_queue_len--;
if (tsq->owner & XMPP_QUEUE_USER)
conn->send_queue_user_len--;
- strophe_free(ctx, tsq);
+ if (!(tsq->owner & XMPP_QUEUE_SM) && conn->sm_state->sm_enabled) {
+ tsq->sm_h = conn->sm_state->sm_sent_nr;
+ conn->sm_state->sm_sent_nr++;
+ strophe_debug_verbose(1, ctx, "xmpp", "SM_Q_MOVE: %p", tsq);
+ add_queue_back(&conn->sm_state->sm_queue, tsq);
+ tsq = NULL;
+ }
+ if (tsq) {
+ strophe_debug_verbose(2, ctx, "xmpp", "Q_FREE: %p", tsq);
+ strophe_debug_verbose(3, ctx, "conn", "Q_CONTENT: %s",
+ tsq->data);
+ strophe_free(ctx, tsq->data);
+ strophe_free(ctx, tsq);
+ }
/* pop the top item */
conn->send_queue_head = sq;
diff --git a/src/tls.c b/src/tls.c
index 5d78445..17d85dd 100644
--- a/src/tls.c
+++ b/src/tls.c
@@ -103,7 +103,7 @@ const char *xmpp_tlscert_get_dnsname(const xmpp_tlscert_t *cert, size_t n)
const char *xmpp_tlscert_get_string(const xmpp_tlscert_t *cert,
xmpp_cert_element_t elmnt)
{
- if (elmnt >= XMPP_CERT_ELEMENT_MAX)
+ if (elmnt < 0 || elmnt >= XMPP_CERT_ELEMENT_MAX)
return NULL;
return cert->elements[elmnt];
}
@@ -132,7 +132,7 @@ const char *xmpp_tlscert_get_description(xmpp_cert_element_t elmnt)
"Fingerprint SHA-1",
"Fingerprint SHA-256",
};
- if (elmnt >= XMPP_CERT_ELEMENT_MAX)
+ if (elmnt < 0 || elmnt >= XMPP_CERT_ELEMENT_MAX)
return NULL;
return descriptions[elmnt];
}
diff --git a/src/tls_openssl.c b/src/tls_openssl.c
index 9fcfb91..4a52947 100644
--- a/src/tls_openssl.c
+++ b/src/tls_openssl.c
@@ -883,9 +883,13 @@ static const char *_tls_error_str(int error, const char **tbl, size_t tbl_size)
static void _tls_set_error(tls_t *tls, int error)
{
if (error != 0 && !tls_is_recoverable(error)) {
- strophe_debug(tls->ctx, "tls", "error=%s(%d) errno=%d",
- TLS_ERROR_STR(error, tls_errors), error, errno);
+ strophe_debug(tls->ctx, "tls", "error=%s(%d) errno=%d lasterror=%d",
+ TLS_ERROR_STR(error, tls_errors), error, errno,
+ tls->lasterror);
_tls_log_error(tls->ctx);
+ } else if (tls->lasterror && tls->lasterror != error) {
+ strophe_debug_verbose(1, tls->ctx, "tls", "overwrite lasterror=%d",
+ tls->lasterror);
}
tls->lasterror = error;
}
diff --git a/src/util.c b/src/util.c
index e3d1794..efde98a 100644
--- a/src/util.c
+++ b/src/util.c
@@ -15,6 +15,7 @@
#include
#include
+#include
#ifdef _WIN32
#include
@@ -161,6 +162,13 @@ void disconnect_mem_error(xmpp_conn_t *conn)
xmpp_disconnect(conn);
}
+int string_to_ul(const char *s, unsigned long *ul)
+{
+ char *endptr;
+ *ul = strtoul(s, &endptr, 10);
+ return *endptr != '\0';
+}
+
void hex_encode(char *writebuf, void *readbuf, size_t len)
{
size_t i;
diff --git a/src/util.h b/src/util.h
index 3fba5e8..f3d7750 100644
--- a/src/util.h
+++ b/src/util.h
@@ -33,6 +33,7 @@ uint64_t time_stamp(void);
uint64_t time_elapsed(uint64_t t1, uint64_t t2);
/* misc functions */
+int string_to_ul(const char *s, unsigned long *ul);
void hex_encode(char *writebuf, void *readbuf, size_t len);
#endif /* __LIBSTROPHE_UTIL_H__ */
diff --git a/strophe.h b/strophe.h
index 20239d3..7e9e7e7 100644
--- a/strophe.h
+++ b/strophe.h
@@ -172,6 +172,7 @@ xmpp_log_t *xmpp_get_default_logger(xmpp_log_level_t level);
/* opaque connection object */
typedef struct _xmpp_conn_t xmpp_conn_t;
typedef struct _xmpp_stanza_t xmpp_stanza_t;
+typedef struct _xmpp_sm_t xmpp_sm_state_t;
/* connection flags */
#define XMPP_CONN_FLAG_DISABLE_TLS (1UL << 0)
@@ -185,6 +186,10 @@ typedef struct _xmpp_stanza_t xmpp_stanza_t;
* Enable legacy authentication support.
*/
#define XMPP_CONN_FLAG_LEGACY_AUTH (1UL << 4)
+/** @def XMPP_CONN_FLAG_DISABLE_SM
+ * Disable Stream-Management XEP-0198.
+ */
+#define XMPP_CONN_FLAG_DISABLE_SM (1UL << 5)
/* connect callback */
typedef enum {
@@ -392,6 +397,11 @@ typedef enum {
char *xmpp_conn_send_queue_drop_element(xmpp_conn_t *conn,
xmpp_queue_element_t which);
+xmpp_sm_state_t *xmpp_conn_get_sm_state(xmpp_conn_t *conn);
+int xmpp_conn_set_sm_state(xmpp_conn_t *conn, xmpp_sm_state_t *sm_state);
+
+void xmpp_free_sm_state(xmpp_sm_state_t *sm_state);
+
int xmpp_connect_client(xmpp_conn_t *conn,
const char *altdomain,
unsigned short altport,
diff --git a/tests/test_send_queue.c b/tests/test_send_queue.c
index af10b86..83897f8 100644
--- a/tests/test_send_queue.c
+++ b/tests/test_send_queue.c
@@ -24,6 +24,7 @@ int main()
xmpp_conn_t *conn;
xmpp_log_t *log;
xmpp_conn_state_t state;
+ xmpp_sm_state_t *sm_state;
char *ret;
unsigned int n;
@@ -32,6 +33,11 @@ int main()
log = xmpp_get_default_logger(XMPP_LEVEL_DEBUG);
ctx = xmpp_ctx_new(NULL, log);
conn = xmpp_conn_new(ctx);
+ sm_state = strophe_alloc(ctx, sizeof(*sm_state));
+ memset(sm_state, 0, sizeof(*sm_state));
+ sm_state->ctx = ctx;
+
+ xmpp_conn_set_sm_state(conn, sm_state);
ENSURE_EQ(xmpp_conn_send_queue_len(conn), 0);
@@ -50,7 +56,7 @@ int main()
xmpp_send_raw(conn, "baan", 4);
ENSURE_EQ(xmpp_conn_send_queue_len(conn), 4);
- conn->send_queue_head->written = 1;
+ conn->send_queue_head->wip = 1;
ENSURE_EQ(xmpp_conn_send_queue_len(conn), 3);
ret = xmpp_conn_send_queue_drop_element(conn, XMPP_QUEUE_OLDEST);
@@ -58,7 +64,7 @@ int main()
xmpp_free(ctx, ret);
ENSURE_EQ(xmpp_conn_send_queue_len(conn), 2);
- conn->send_queue_head->written = 0;
+ conn->send_queue_head->wip = 0;
ENSURE_EQ(xmpp_conn_send_queue_len(conn), 3);
ret = xmpp_conn_send_queue_drop_element(conn, XMPP_QUEUE_OLDEST);