diff --git a/src/builder/b-artifacts.c b/src/builder/b-artifacts.c
index a51f299..4f3f70b 100644
--- a/src/builder/b-artifacts.c
+++ b/src/builder/b-artifacts.c
@@ -48,6 +48,7 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf,
size_t *len, int *flags)
{
sai_artifact_t *ap = (sai_artifact_t *)userobj;
+ lws_ss_state_return_t r;
lws_struct_serialize_t *js;
size_t w;
int n;
@@ -68,7 +69,9 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf,
*len = w;
ap->sent_json = 1;
- lws_ss_request_tx(ap->ss);
+ r = lws_ss_request_tx(ap->ss);
+ if (r)
+ return r;
lwsl_notice("%s: sent JSON %s\n", __func__, (const char *)buf);
return LWSSSSRET_OK;
@@ -110,9 +113,7 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf,
}
/* even if we finished, we want to come back to close */
- lws_ss_request_tx(ap->ss);
-
- return LWSSSSRET_OK;
+ return lws_ss_request_tx(ap->ss);
}
static lws_ss_state_return_t
@@ -160,8 +161,7 @@ saib_artifact_state(void *userobj, void *sh, lws_ss_constate_t state,
case LWSSSCS_CONNECTED:
// lwsl_user("%s: CONNECTED: %p\n", __func__, ap->ss);
- lws_ss_request_tx_len(ap->ss, (unsigned long)ap->len);
- break;
+ return lws_ss_request_tx_len(ap->ss, (unsigned long)ap->len);
case LWSSSCS_TIMEOUT:
lwsl_info("%s: timeout\n", __func__);
diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c
index e3dc155..12dbda4 100644
--- a/src/builder/b-comms.c
+++ b/src/builder/b-comms.c
@@ -213,6 +213,7 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
struct sai_plat *sp = NULL;
lws_struct_serialize_t *js;
lws_dll2_t *star, *walk;
+ lws_ss_state_return_t r;
struct sai_nspawn *ns;
size_t w = 0;
int n = 0;
@@ -253,7 +254,9 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
lws_dll2_remove(&rej->list);
free(rej);
- lws_ss_request_tx(spm->ss);
+ r = lws_ss_request_tx(spm->ss);
+ if (r)
+ return r;
goto sendify;
}
@@ -278,7 +281,9 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
lwsl_notice("%s: forwarding to server %.*s\n", __func__,
(int)(*len), (const char *)buf);
- lws_ss_request_tx(spm->ss);
+ r = lws_ss_request_tx(spm->ss);
+ if (r)
+ return r;
goto sendify;
}
@@ -311,9 +316,9 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
*flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM;
if (spm->logs_in_flight)
- lws_ss_request_tx(spm->ss);
+ return lws_ss_request_tx(spm->ss);
- return 0;
+ return LWSSSSRET_OK;
}
/*
@@ -460,13 +465,16 @@ sendify:
*flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM;
*len = (unsigned int)n;
- if (spm->phase != PHASE_IDLE || spm->logs_in_flight)
- lws_ss_request_tx(spm->ss);
+ if (spm->phase != PHASE_IDLE || spm->logs_in_flight) {
+ r = lws_ss_request_tx(spm->ss);
+ if (r)
+ return r;
+ }
if (!n)
return 1;
- return 0;
+ return LWSSSSRET_OK;
}
static int
@@ -618,8 +626,7 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
case LWSSSCS_CONNECTED:
lwsl_user("%s: CONNECTED: %p\n", __func__, spm->ss);
spm->phase = PHASE_START_ATTACH;
- lws_ss_request_tx(spm->ss);
- break;
+ return lws_ss_request_tx(spm->ss);
case LWSSSCS_DISCONNECTED:
/*
@@ -633,8 +640,7 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
case LWSSSCS_ALL_RETRIES_FAILED:
lwsl_user("%s: LWSSSCS_ALL_RETRIES_FAILED\n", __func__);
- lws_ss_request_tx(spm->ss);
- break;
+ return lws_ss_request_tx(spm->ss);
case LWSSSCS_QOS_ACK_REMOTE:
lwsl_notice("%s: LWSSSCS_QOS_ACK_REMOTE\n", __func__);
@@ -644,7 +650,7 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
break;
}
- return 0;
+ return LWSSSSRET_OK;
}
const lws_ss_info_t ssi_sai_builder = {
diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c
index 1b771f4..af510d8 100644
--- a/src/builder/b-nspawn.c
+++ b/src/builder/b-nspawn.c
@@ -121,8 +121,7 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason,
if (!saib_log_chunk_create(ns, buf, len, lws_spawn_get_stdfd(wsi)))
return -1;
- lws_ss_request_tx(ns->spm->ss);
- break;
+ return lws_ss_request_tx(ns->spm->ss) ? -1 : 0;
default:
break;
diff --git a/src/builder/b-refproxy.c b/src/builder/b-refproxy.c
index e7adc15..097a5f0 100644
--- a/src/builder/b-refproxy.c
+++ b/src/builder/b-refproxy.c
@@ -92,9 +92,8 @@ saib_queue_yield_message(struct sai_plat_server *spm, const char *c, size_t len)
lwsl_notice("%s: %s\n", __func__, m->msg);
lws_dll2_add_tail(&m->list, &spm->resource_req_list);
- lws_ss_request_tx(spm->ss);
- return 0;
+ return lws_ss_request_tx(spm->ss) ? -1 : 0;
}
int
@@ -226,8 +225,7 @@ callback_resproxy(struct lws *wsi, enum lws_callback_reasons reason,
lws_dll2_add_tail(&pss->list, &spm->resource_pss_list);
/* try to schedule a write */
- lws_ss_request_tx(spm->ss);
- break;
+ return lws_ss_request_tx(spm->ss) ? -1 : 0;
case LWS_CALLBACK_RAW_WRITEABLE:
if (pss->response) {
diff --git a/src/builder/b-task.c b/src/builder/b-task.c
index a629084..537577e 100644
--- a/src/builder/b-task.c
+++ b/src/builder/b-task.c
@@ -77,7 +77,7 @@ saib_set_ns_state(struct sai_nspawn *ns, int state)
}
if (ns->spm && ns->spm->ss)
- lws_ss_request_tx(ns->spm->ss);
+ return lws_ss_request_tx(ns->spm->ss) ? -1 : 0;
return 0;
}
@@ -125,9 +125,8 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm,
rej->ongoing = sp->ongoing;
lws_dll2_add_tail(&rej->list, &spm->rejection_list);
- lws_ss_request_tx(spm->ss);
- return 0;
+ return lws_ss_request_tx(spm->ss) ? -1 : 0;
}
void
@@ -143,7 +142,8 @@ saib_task_destroy(struct sai_nspawn *ns)
*/
if (ns->spm) {
ns->spm->phase = PHASE_START_ATTACH;
- lws_ss_request_tx(ns->spm->ss);
+ if (lws_ss_request_tx(ns->spm->ss))
+ return;
}
if (ns->tp) {
@@ -268,9 +268,8 @@ artifact_glob_cb(void *data, const char *path)
if (lws_ss_set_metadata(h, "url", ns->spm->url, strlen(ns->spm->url)))
lwsl_warn("%s: unable to set metadata\n", __func__);
- lws_ss_client_connect(h);
- return 0;
+ return lws_ss_client_connect(h) ? -1 : 0;
}
/*
diff --git a/src/common/ss-client-logproxy.c b/src/common/ss-client-logproxy.c
index 36473e0..166e6e7 100644
--- a/src/common/ss-client-logproxy.c
+++ b/src/common/ss-client-logproxy.c
@@ -67,7 +67,10 @@ saicom_lp_ss_from_env(struct lws_context *context, const char *env_name)
if (lws_ss_set_metadata(h, "sockpath", e, strlen(e)))
lwsl_warn("%s: metadata set failed\n", __func__);
- lws_ss_client_connect(h);
+ if (lws_ss_client_connect(h)) {
+ lws_ss_destroy(&h);
+ return NULL;
+ }
return h;
}
@@ -85,9 +88,8 @@ saicom_lp_add(struct lws_ss_handle *h, const char *buf, size_t len)
log->len = len;
lws_dll2_add_tail(&log->list, &lp->logs);
- lws_ss_request_tx(h);
- return 0;
+ return lws_ss_request_tx(h);
}
void
@@ -135,6 +137,7 @@ saicom_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
{
said_logproxy_t *lp = (said_logproxy_t *)userobj;
said_log_t *log = lws_container_of(lp->logs.head, said_log_t, list);
+ lws_ss_state_return_t r;
if (!lp->logs.head)
return 1; /* nothing to send */
@@ -154,9 +157,9 @@ saicom_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
}
if (lp->logs.head) {
- lws_ss_request_tx(lp->ss);
-
- return 0;
+ r = lws_ss_request_tx(lp->ss);
+ if (r)
+ return r;
}
check_drained();
@@ -178,8 +181,7 @@ saicom_lp_state(void *userobj, void *sh, lws_ss_constate_t state,
break;
case LWSSSCS_CONNECTED:
- lws_ss_request_tx(lp->ss);
- break;
+ return lws_ss_request_tx(lp->ss);
case LWSSSCS_DISCONNECTED:
break;
diff --git a/src/resource/r-comms.c b/src/resource/r-comms.c
index 2d0faac..f78279a 100644
--- a/src/resource/r-comms.c
+++ b/src/resource/r-comms.c
@@ -82,8 +82,7 @@ sair_lp_state(void *userobj, void *sh, lws_ss_constate_t state,
switch (state) {
case LWSSSCS_CONNECTED:
- lws_ss_request_tx(lp->ss);
- break;
+ return lws_ss_request_tx(lp->ss);
case LWSSSCS_ALL_RETRIES_FAILED:
case LWSSSCS_DISCONNECTED:
@@ -132,7 +131,10 @@ sair_ss_from_env(struct lws_context *context, const char *env_name)
if (lws_ss_set_metadata(h, "sockpath", e, strlen(e)))
lwsl_err("%s: unable to set metadata\n", __func__);
- lws_ss_client_connect(h);
+ if (lws_ss_client_connect(h)) {
+ lws_ss_destroy(&h);
+ return NULL;
+ }
return h;
}
diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c
index caf1d25..5d39fd4 100644
--- a/src/server/s-websrv.c
+++ b/src/server/s-websrv.c
@@ -104,7 +104,8 @@ sais_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len)
n = lws_buflist_append_segment(&m->bltx, buf, len);
lwsl_notice("%s: appened h %p: %d\n", __func__, h, n);
- lws_ss_request_tx(h);
+ if (lws_ss_request_tx(h))
+ return 1;
return n < 0;
}
@@ -120,10 +121,14 @@ _sais_websrv_broadcast(struct lws_ss_handle *h, void *arg)
websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h);
sais_websrv_broadcast_t *a = (sais_websrv_broadcast_t *)arg;
- if (lws_buflist_append_segment(&m->bltx, a->buf, a->len) >= 0)
- lws_ss_request_tx(h);
- else
+ if (lws_buflist_append_segment(&m->bltx, a->buf, a->len) < 0) {
lwsl_warn("%s: buflist append fail\n", __func__);
+
+ return;
+ }
+
+ if (lws_ss_request_tx(h))
+ lwsl_ss_warn(h, "tx req fail");
}
void
@@ -207,10 +212,14 @@ _sais_taskchange(struct lws_ss_handle *h, void *_arg)
"\"event_hash\":\"%s\", \"state\":%d}",
arg->uid, arg->state);
- if (lws_buflist_append_segment(&m->bltx, (uint8_t *)tc, (unsigned int)n) >= 0)
- lws_ss_request_tx(h);
- else
+ if (lws_buflist_append_segment(&m->bltx, (uint8_t *)tc, (unsigned int)n) < 0) {
lwsl_warn("%s: buflist append failed\n", __func__);
+
+ return;
+ }
+
+ if (lws_ss_request_tx(h))
+ lwsl_ss_warn(h, "tx req fail");
}
void
@@ -221,7 +230,7 @@ sais_taskchange(struct lws_ss_handle *hsrv, const char *task_uuid, int state)
lws_ss_server_foreach_client(hsrv, _sais_taskchange, (void *)&arg);
}
-void
+static void
_sais_eventchange(struct lws_ss_handle *h, void *_arg)
{
websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h);
@@ -233,10 +242,13 @@ _sais_eventchange(struct lws_ss_handle *h, void *_arg)
"\"event_hash\":\"%s\", \"state\":%d}",
arg->uid, arg->state);
- if (lws_buflist_append_segment(&m->bltx, (uint8_t *)tc, (unsigned int)n) >= 0)
- lws_ss_request_tx(h);
- else
+ if (lws_buflist_append_segment(&m->bltx, (uint8_t *)tc, (unsigned int)n) < 0) {
lwsl_warn("%s: buflist append failed\n", __func__);
+ return;
+ }
+
+ if (lws_ss_request_tx(h))
+ lwsl_ss_warn(h, "req fail");
}
void
@@ -466,7 +478,7 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf,
*len = (size_t)used;
if (m->bltx)
- lws_ss_request_tx(m->ss);
+ return lws_ss_request_tx(m->ss);
return 0;
}
@@ -485,8 +497,8 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state,
lws_buflist_destroy_all_segments(&m->bltx);
break;
case LWSSSCS_CREATING:
- lws_ss_request_tx(m->ss);
- break;
+ return lws_ss_request_tx(m->ss);
+
case LWSSSCS_CONNECTED:
sais_list_builders(m->vhd);
break;
diff --git a/src/web/w-comms.c b/src/web/w-comms.c
index e42ea09..e8c50f2 100644
--- a/src/web/w-comms.c
+++ b/src/web/w-comms.c
@@ -513,9 +513,8 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user,
if (lws_ss_set_metadata(vhd->h_ss_websrv, "sockpath",
"@com.warmcat.sai-websrv", 23))
lwsl_warn("%s: unable to set metadata\n", __func__);
- lws_ss_client_connect(vhd->h_ss_websrv);
- break;
+ return lws_ss_client_connect(vhd->h_ss_websrv) ? -1 : 0;
case LWS_CALLBACK_PROTOCOL_DESTROY:
saiw_event_db_close_all_now(vhd);
diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c
index a3d3257..8ab9359 100644
--- a/src/web/w-websrv.c
+++ b/src/web/w-websrv.c
@@ -182,7 +182,7 @@ saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
*len = (size_t)used;
if (m->bltx)
- lws_ss_request_tx(m->ss);
+ return lws_ss_request_tx(m->ss);
return 0;
}
@@ -203,8 +203,7 @@ saiw_lp_state(void *userobj, void *sh, lws_ss_constate_t state,
case LWSSSCS_CONNECTED:
lwsl_notice("%s: connected to websrv uds\n", __func__);
- lws_ss_request_tx(m->ss);
- break;
+ return lws_ss_request_tx(m->ss);
case LWSSSCS_DISCONNECTED:
lws_buflist_destroy_all_segments(&m->bltx);
@@ -212,8 +211,7 @@ saiw_lp_state(void *userobj, void *sh, lws_ss_constate_t state,
break;
case LWSSSCS_ALL_RETRIES_FAILED:
- lws_ss_client_connect(m->ss);
- break;
+ return lws_ss_client_connect(m->ss);
case LWSSSCS_QOS_ACK_REMOTE:
break;
@@ -247,7 +245,5 @@ saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len)
if (lws_buflist_append_segment(&m->bltx, buf, len) < 0)
return 1;
- lws_ss_request_tx(h);
-
- return 0;
+ return !!lws_ss_request_tx(h);
}