diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c
index 6ee11e7..52298f0 100644
--- a/src/builder/b-nspawn.c
+++ b/src/builder/b-nspawn.c
@@ -51,6 +51,7 @@ saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel)
if (!ns->task)
return 0;
+
n = lws_snprintf(lj + LWS_PRE, sizeof(lj) - LWS_PRE,
"{\"schema\":\"com-warmcat-sai-logs\","
"\"task_uuid\":\"%s\", \"timestamp\": %llu,"
@@ -100,6 +101,12 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason,
switch (reason) {
case LWS_CALLBACK_RAW_CLOSE_FILE:
+ {
+ int ch = lws_spawn_get_stdfd(wsi);
+ if (ch == 0) ch = 1;
+ if (op && op->ns && ch < 3)
+ op->ns->stdwsi[ch] = NULL;
+ }
if (op && op->lsp) {
if (lws_spawn_stdwsi_closed(op->lsp, wsi) &&
ns->reap_cb_called) {
@@ -142,8 +149,32 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason,
int ch = lws_spawn_get_stdfd(wsi);
if (ch == 0)
ch = 1;
+
+ if (ch < 3)
+ op->ns->stdwsi[ch] = wsi;
+
if (saib_log_chunk_create(op->ns, buf, len, ch))
return -1;
+
+ if (lws_buflist2_total_len(&op->ns->spm->bl_to_srv) > (LWS_BUFLIST_OOM_LIMIT - (256 * 1024))) {
+ /* buflist is getting full, backpressure ALL active stdwsi for this connection */
+ lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, builder.sai_plat_owner.head) {
+ sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list);
+ lws_start_foreach_dll_safe(struct lws_dll2 *, d2, d3, sp->nspawn_owner.head) {
+ struct sai_nspawn *ns = lws_container_of(d2, struct sai_nspawn, list);
+ if (ns->spm == op->ns->spm) {
+ for (int i = 0; i < 3; i++) {
+ if (!ns->stdwsi_paused[i] && ns->stdwsi[i]) {
+ ns->stdwsi_paused[i] = 1;
+ lws_rx_flow_control(ns->stdwsi[i], 0); /* 0 disables RX */
+ lwsl_notice("%s: Backpressure applied to ch %d (tot %zu)\n",
+ __func__, i, lws_buflist2_total_len(&op->ns->spm->bl_to_srv));
+ }
+ }
+ }
+ } lws_end_foreach_dll_safe(d2, d3);
+ } lws_end_foreach_dll_safe(d, d1);
+ }
}
return lws_ss_request_tx(op->ns->spm->ss) ? -1 : 0;
diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c
index 0c66feb..93c4c86 100644
--- a/src/builder/b-sai.c
+++ b/src/builder/b-sai.c
@@ -520,7 +520,7 @@ crash_handler(int signum)
int
saib_app_run(int argc, const char **argv)
{
- int logs = 1039 | LLL_USER | LLL_ERR | LLL_WARN | LLL_NOTICE;
+ int logs = LLL_USER | LLL_ERR | LLL_WARN | LLL_NOTICE;
struct lws_context_creation_info info;
#if defined(WIN32)
char temp[256], stg_config_dir[256];
diff --git a/src/builder/b-ws-server.c b/src/builder/b-ws-server.c
index 210ca28..3a1acd8 100644
--- a/src/builder/b-ws-server.c
+++ b/src/builder/b-ws-server.c
@@ -82,9 +82,12 @@ saib_srv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int s
// lwsl_ss_notice(h, "Queuing builder -> sai-server");
// lwsl_hexdump_notice(buf, len);
- if (lws_buflist_append_segment(&spm->bl_to_srv, (uint8_t *)buf - sizeof(int),
- len + sizeof(int)) < 0)
- lwsl_ss_err(h, "failed to append"); /* still ask to drain */
+ if (lws_buflist2_append_segment(&spm->bl_to_srv, (uint8_t *)buf - sizeof(int),
+ len + sizeof(int)) < 0) {
+ lwsl_err("%s: failed to append\n", __func__);
+ spm->tx_corrupted = 1; /* Mark the stream as permanently corrupted */
+ return -1;
+ }
if (lws_ss_request_tx(h))
lwsl_ss_err(h, "failed to request tx");
@@ -118,6 +121,10 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h,
break;
case LSJS_RESULT_ERROR:
lwsl_warn("%s: serialization failed\n", __func__);
+ {
+ struct sai_plat_server *spm = (struct sai_plat_server *)lws_ss_to_user_object(h);
+ spm->tx_corrupted = 1;
+ }
return -1;
}
@@ -316,20 +323,26 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
int *flags)
{
struct sai_plat_server *spm = (struct sai_plat_server *)userobj;
- unsigned int *pi = (unsigned int *)lws_buflist_get_frag_start_or_NULL(&spm->bl_to_srv);
+ unsigned int *pi = (unsigned int *)lws_buflist2_get_frag_start_or_NULL(&spm->bl_to_srv);
char som, som1, eom, final = 1;
size_t fsl, used;
- if (!spm->bl_to_srv)
+ if (spm->tx_corrupted) {
+ lwsl_err("%s: tx_corrupted, dropping connection\n", __func__);
+ return LWSSSSRET_DISCONNECT_ME;
+ }
+
+ if (!spm->bl_to_srv.owner.head)
return LWSSSSRET_TX_DONT_SEND;
- fsl = lws_buflist_next_segment_len(&spm->bl_to_srv, NULL);
+ fsl = lws_buflist2_next_segment_len(&spm->bl_to_srv, NULL);
+
+ lws_buflist2_fragment_use(&spm->bl_to_srv, NULL, 0, &som, &eom);
- lws_buflist_fragment_use(&spm->bl_to_srv, NULL, 0, &som, &eom);
if (som) {
spm->tx_flags = *pi;
fsl -= sizeof(int);
- lws_buflist_fragment_use(&spm->bl_to_srv, buf, sizeof(int), &som1, &eom);
+ lws_buflist2_fragment_use(&spm->bl_to_srv, buf, sizeof(int), &som1, &eom);
}
if (!(spm->tx_flags & LWSSS_FLAG_SOM))
som = 0;
@@ -337,7 +350,7 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
if (fsl == 0)
used = 0;
else
- used = (size_t)lws_buflist_fragment_use(&spm->bl_to_srv, (uint8_t *)buf, *len, &som1, &eom);
+ used = (size_t)lws_buflist2_fragment_use(&spm->bl_to_srv, (uint8_t *)buf, *len, &som1, &eom);
if (used < fsl || !(spm->tx_flags & LWSSS_FLAG_EOM))
final = 0;
@@ -360,13 +373,27 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
if (*flags & LWSSS_FLAG_EOM)
spm->inside_msg = 0;
-// lwsl_ss_notice(spm->ss, "Sending %d builder->srv: ssflags %d", (int)*len, (int)*flags);
-// lwsl_hexdump_notice(buf, *len);
-
- if (spm->bl_to_srv)
+ if (spm->bl_to_srv.owner.head)
return lws_ss_request_tx(spm->ss);
- return 0;
+ /* buflist is empty, unpause any backpressured stdwsi */
+ lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, builder.sai_plat_owner.head) {
+ sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list);
+ lws_start_foreach_dll_safe(struct lws_dll2 *, d2, d3, sp->nspawn_owner.head) {
+ struct sai_nspawn *ns = lws_container_of(d2, struct sai_nspawn, list);
+ if (ns->spm == spm) {
+ for (int i = 0; i < 3; i++) {
+ if (ns->stdwsi_paused[i] && ns->stdwsi[i]) {
+ lwsl_notice("%s: Unpausing ch %d\n", __func__, i);
+ ns->stdwsi_paused[i] = 0;
+ lws_rx_flow_control(ns->stdwsi[i], 1 | LWS_RXFLOW_REASON_USER_BOOL);
+ }
+ }
+ }
+ } lws_end_foreach_dll_safe(d2, d3);
+ } lws_end_foreach_dll_safe(d, d1);
+
+ return LWSSSSRET_OK;
}
static int
@@ -644,6 +671,13 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
case LWSSSCS_CONNECTED:
lwsl_ss_user(spm->ss, "CONNECTED");
+
+ /* Reset corruption flags and discard broken buflist on reconnect */
+ spm->tx_corrupted = 0;
+ spm->inside_msg = 0;
+ spm->last_msg_start[0] = '\0';
+ lws_buflist2_destroy_all_segments(&spm->bl_to_srv);
+
/* Initialize the load report SUL timer for this server connection */
lws_sul_schedule(builder.context, 0, &spm->sul_load_report,
saib_sul_load_report_cb, 1);
diff --git a/src/common/include/private.h b/src/common/include/private.h
index 2bd1d07..4d76d91 100644
--- a/src/common/include/private.h
+++ b/src/common/include/private.h
@@ -251,6 +251,9 @@ struct sai_nspawn {
struct saib_opaque_spawn *op;
sai_task_t *task;
+ struct lws *stdwsi[3];
+ uint8_t stdwsi_paused[3];
+
#if defined(LWS_WITH_SPAWN)
lws_spawn_resource_us_t res;
#endif
@@ -471,7 +474,7 @@ typedef struct sai_plat_server {
lws_dll2_owner_t resource_pss_list; /* so we can find the cookie */
- struct lws_buflist *bl_to_srv;
+ struct lws_buflist2_owner bl_to_srv;
char resproxy_path[128];
@@ -497,6 +500,7 @@ typedef struct sai_plat_server {
char last_msg_start[128];
uint8_t inside_msg;
+ uint8_t tx_corrupted;
} sai_plat_server_t;
struct sai_env {