Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / assets / solaris-11.svg
Author[]Andy Green <andy@warmcat.com> 2025-09-01 05:06 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-09-02 15:10 UTC
Tree428ebed2f94bf460604c42ad569146350ec71795   Raw Patch
 
sai-web: fix incorrect FIN handling possible
sai-web: fix incorrect FIN handling possible
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;
Page fetched 0s ago, creation time: 4ms (vhost etag hits: 0%, cache hits: 0%)