diff --git a/src/web/w-private.h b/src/web/w-private.h
index 3380522..6cd0817 100644
--- a/src/web/w-private.h
+++ b/src/web/w-private.h
@@ -320,7 +320,7 @@ saiw_sched_destroy(struct lws_dll2 *d, void *user);
void
-saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int min_api_version, int flags);
+saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int min_api_version, enum lws_write_protocol flags);
void
saiw_browser_state_changed(struct pss *pss, int established);
diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c
index 1e16c29..ead565f 100644
--- a/src/web/w-websrv.c
+++ b/src/web/w-websrv.c
@@ -89,10 +89,6 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
goto cleanup_and_disconnect;
}
- /*
- * This is the key: if the message is not yet complete, just return
- * and wait for the next fragment. Don't process anything yet.
- */
if (n == LEJP_CONTINUE) {
/*
* Also forward this fragment to browsers if the message is for them.
@@ -102,7 +98,11 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
switch (m->a.top_schema_index) {
case SAIS_WS_WEBSRV_RX_LOADREPORT:
saiw_ws_broadcast_raw(vhd, buf, len, 0,
- lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM));
+ lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, 0));
+ break;
+ case SAIS_WS_WEBSRV_RX_TASKACTIVITY:
+ saiw_ws_broadcast_raw(vhd, buf, len, 0,
+ lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, 0));
break;
}
@@ -174,7 +174,7 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
case SAIS_WS_WEBSRV_RX_LOADREPORT:
/* Forward the final fragment of the load report */
- saiw_ws_broadcast_raw(vhd, buf, len, 2,
+ saiw_ws_broadcast_raw(vhd, buf, len, 0,
lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM));
break;
case SAIS_WS_WEBSRV_RX_TASKACTIVITY:
@@ -203,19 +203,31 @@ saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
int *flags)
{
saiw_websrv_t *m = (saiw_websrv_t *)userobj;
- char som, eom;
+ char som, eom, final = 1;
+ size_t fsl;
int used;
if (!m->bltx)
return LWSSSSRET_TX_DONT_SEND;
+ /*
+ * We can only issue *len at a time.
+ */
+
+ fsl = lws_buflist_next_segment_len(&m->bltx, NULL);
+
+ if (fsl > *len)
+ final = 0;
+
used = lws_buflist_fragment_use(&m->bltx, buf, *len, &som, &eom);
if (!used)
return LWSSSSRET_TX_DONT_SEND;
- *flags = (som ? LWSSS_FLAG_SOM : 0) | (eom ? LWSSS_FLAG_EOM : 0);
+ *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0);
*len = (size_t)used;
+ lwsl_ss_notice(m->ss, "TX issuing %d bytes, ss flags %d\n", (int)used, *flags);
+
if (m->bltx)
return lws_ss_request_tx(m->ss);
diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c
index 68775d1..25c3c9a 100644
--- a/src/web/w-ws-browser.c
+++ b/src/web/w-ws-browser.c
@@ -33,28 +33,27 @@
/*
* This allows other parts of sai-web to queue a raw buffer to be sent to
* all connected browsers, eg, for load reports.
+ *
+ * The flags are lws_write() flags.
*/
void
-saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int api_ver_min, int flags)
+saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int api_ver_min, enum lws_write_protocol flags)
{
int eff = 0;
- // lwsl_err("%s: sai-web broadcasting to browsers\n", __func__);
- // lwsl_hexdump_err(buf, len);
lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) {
struct pss *pss = lws_container_of(p, struct pss, same);
- int *pi = (int *)((const char *)buf -sizeof(int));
+ int *pi = (int *)((const char *)buf - sizeof(int));
if (pss->js_api_version >= api_ver_min) {
eff++;
- *pi = flags;
+ *pi = (int)flags;
if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) > 0)
lws_callback_on_writable(pss->wsi);
}
- } lws_end_foreach_dll(p);
- // lwsl_notice("%s: broadcast to %d / %d browsers\n", __func__, eff, (int)vhd->browsers.count);
+ } lws_end_foreach_dll(p);
}
extern const lws_struct_map_t lsm_load_report_members[7];
@@ -790,8 +789,11 @@ again:
* Stay in this state if we're in the middle of a
* multi-fragment message
*/
- if (lws_ws_sending_multifragment(pss->wsi))
+ if (lws_ws_sending_multifragment(pss->wsi)) {
+ lws_callback_on_writable(pss->wsi);
+
return 0;
+ }
/* fallthru */
@@ -804,27 +806,40 @@ again:
*/
if (pss->raw_tx) {
- char som, eom, rb[4096];
- int used, *pi = (int *)rb;
+ /*
+ * Notice we are getting the stored flags from the START of the fragment each time.
+ * that means we can still see the right flags stored with the fragment, even if we
+ * have partially used the buflist frag and are partway through it.
+ *
+ * Ergo, only something to skip if we are at som=1. And also notice that although
+ * *pi will be right, after the lws_buflist..._use() api, what it points to has been
+ * destroyed. So we also dereference *pi into depi for use below.
+ */
+ int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&pss->raw_tx), depi = *pi;
+ char som, eom, rb[1200];
+ int used, final = 1;
+ size_t fsl;
- used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb,
- sizeof(rb), &som, &eom);
+ fsl = lws_buflist_next_segment_len(&pss->raw_tx, NULL);
+
+ /* this is the only buflist user on pss->raw_tx */
+ used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, sizeof(rb), &som, &eom);
if (!used)
return 0;
+ if (used < (int)fsl || (depi & LWS_WRITE_NO_FIN))
+ final = 0;
- // lwsl_wsi_notice(pss->wsi, "writing %d bytes flags 0x%x: '%.*s'",
- // (int)(used - (int)sizeof(int)), (int)*pi,
- // (int)(used - (int)sizeof(int)), rb + sizeof(int));
-
- if (lws_write(pss->wsi, (uint8_t *)rb + sizeof(int),
- (size_t)used - sizeof(int),
- (enum lws_write_protocol)*pi) < 0) {
+ if (lws_write(pss->wsi, (uint8_t *)rb + ((size_t)som * sizeof(int)),
+ (size_t)used - ((size_t)som * sizeof(int)),
+ lws_write_ws_flags((((enum lws_write_protocol)depi) & 0xf), som, final)
+ ) < 0) {
lwsl_wsi_err(pss->wsi, "attempt to write %d failed", (int)used - (int)sizeof(int));
return -1;
}
- lws_callback_on_writable(pss->wsi);
+ if (lws_buflist_next_segment_len(&pss->raw_tx, NULL))
+ lws_callback_on_writable(pss->wsi);
if (!lws_ws_sending_multifragment(pss->wsi))
pss->send_state = WSS_IDLE1;
@@ -842,6 +857,8 @@ again:
!sch)
return 0;
+ /* switch to the pending sch */
+
pss->toggle_favour_sch = 0;
pss->send_state = sch->action;
goto again;