Author: Andy Green Date: Tue Sep 02 09:40:35 2025 +0100 js: taskreset try to trigger logs subscription diff --git a/assets/sai.js b/assets/sai.js index 9eec0d8..74a06f7 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -1575,6 +1575,23 @@ function ws_open_sai() console.log(rs); sai.send(rs); + + /* + * and immediately re-request the task info, so we can get + * the new logs + */ + var tid = san(e.srcElement.id.substring(8)); + var rq = "{\"schema\":" + + "\"com.warmcat.sai.taskinfo\"," + + "\"js_api_version\": " + SAI_JS_API_VERSION + "," + + "\"logs\": 1," + + "\"last_log_ts\":" + last_log_timestamp + "," + + "\"task_hash\":" + + JSON.stringify(tid) + "}"; + + console.log(rq); + sai.send(rq); + document.getElementById("dlogsn").innerHTML = ""; document.getElementById("dlogst").innerHTML = ""; document.getElementById("logs").innerHTML = ""; diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index 19c45f2..907d9e5 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -269,6 +269,8 @@ send_logs: star = NULL; do { + uint32_t tries = builder.sai_plat_owner.count; + if (!spm->last_logging_nspawn) { /* start at the start */ sp = spm->last_logging_platform = lws_container_of( @@ -282,7 +284,7 @@ send_logs: } /* if no more nspawns, try moving to next platform */ - if (!walk) { + while (!walk && tries--) { if (!sp->sai_plat_list.next) /* if no more platforms, wrap around to first */ sp = spm->last_logging_platform = @@ -304,6 +306,7 @@ send_logs: if (walk == star) { lwsl_notice("%s: did not find logs: %d expected\n", __func__, spm->logs_in_flight); + spm->logs_in_flight = 0; return 1; /* nothing to do */ } @@ -613,6 +616,8 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, lws_sul_cancel(&spm->sul_load_report); lws_dll2_foreach_safe(&builder.sai_plat_owner, spm, cleanup_on_ss_disconnect); + if (lws_ss_request_tx(spm->ss)) + lwsl_err("%s: failed to reconnect\n", __func__); break; case LWSSSCS_ALL_RETRIES_FAILED: diff --git a/src/server/s-notification.c b/src/server/s-notification.c index 0d6aeae..cab06c0 100644 --- a/src/server/s-notification.c +++ b/src/server/s-notification.c @@ -267,6 +267,7 @@ sai_saifile_lejp_cb(struct lejp_ctx *ctx, char reason) sn->t.prep[0] = '\0'; sn->t.packages[0] = '\0'; sn->t.cmake[0] = '\0'; + sn->t.cpack[0] = '\0'; sn->t.artifacts[0] = '\0'; sn->explicit_platforms[0] = '\0'; return 0; diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index b09ac29..58ac7d5 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -37,6 +37,7 @@ #include #include #include +#include #include "s-private.h" @@ -159,20 +160,15 @@ reject: return 1; } -int -sais_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len) -{ - websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); - int n; - - n = lws_buflist_append_segment(&m->bltx, buf, len); - - lwsl_notice("%s: appened h %p: %d\n", __func__, h, n); - if (lws_ss_request_tx(h)) - return 1; - - return n < 0; -} +/* + * sais_webserv_broadcast allows us to queue to broadcast a message to all + * sai-web daemons that are connected to us. + * + * The queue is drained by websrvss_ws_tx() below. + * + * These messages are defined to all fit in a single fragment and will + * cause an assertion if they don't. + */ typedef struct { const uint8_t *buf; @@ -194,7 +190,7 @@ _sais_websrv_broadcast(struct lws_ss_handle *h, void *arg) } if (lws_buflist_append_segment(&m->bltx, a->buf, a->len) < 0) { - lwsl_warn("%s: buflist append fail\n", __func__); + lwsl_err("%s: buflist append fail\n", __func__); lws_ss_start_timeout(h, 1); return; @@ -215,6 +211,7 @@ sais_websrv_broadcast(struct lws_ss_handle *hsrv, const char *str, size_t len) lws_ss_server_foreach_client(hsrv, _sais_websrv_broadcast, &a); } + int sais_list_builders(struct vhd *vhd) { @@ -744,6 +741,7 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, int *flags) { websrvss_srv_t *m = (websrvss_srv_t *)userobj; + size_t fsl = lws_buflist_next_segment_len(&m->bltx, NULL); char som, eom; int used; @@ -754,10 +752,20 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, if (!used) return LWSSSSRET_TX_DONT_SEND; + // if (!eom) { + // lwsl_err("%s: eom is not set on bltx buflist!\n", __func__); + // lwsl_hexdump_notice(buf, (size_t)used); + // } + + if ((size_t)used < fsl) { + lwsl_ss_warn(m->ss, "srv->web: clearing eom since used %d < fsl %d\n", (int)used, (int)fsl); + eom = 0; /* because we still be back */ + } + *flags = (som ? LWSSS_FLAG_SOM : 0) | (eom ? LWSSS_FLAG_EOM : 0); *len = (size_t)used; - // lwsl_warn("%s: srv -> web: len %d flags %d\n", __func__, (int)*len, (int)*flags); + lwsl_ss_warn(m->ss, "srv -> web: len %d flags %d", (int)*len, (int)*flags); if (m->bltx) return lws_ss_request_tx(m->ss); diff --git a/src/web/w-comms.c b/src/web/w-comms.c index 4427885..1ddf2cf 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -107,24 +107,6 @@ size_t lsm_schema_json_map_array_size = LWS_ARRAY_SIZE(lsm_schema_json_map); extern const lws_struct_map_t lsm_schema_sq3_map_event[]; -#if 0 -/* - * Let the server know how many browsers are connected, so it can inform - * builders who can then moderate their reporting rate - */ -void -saiw_update_viewer_count(struct vhd *vhd) -{ - char buf[128]; - int n; - - n = lws_snprintf(buf, sizeof(buf), - "{\"schema\":\"com.warmcat.sai.viewercount\",\"count\":%u}", - (unsigned int)vhd->browsers.count); - saiw_websrv_queue_tx(vhd->h_ss_websrv, (uint8_t *)buf, (size_t)n); -} -#endif - /* len is typically 16 (event uuid is 32 chars + NUL) * But eg, task uuid is concatenated 32-char eventid and 32-char taskid */ @@ -1224,7 +1206,8 @@ clean_spa: /* * Browser UI sent us something on websockets */ - if (saiw_ws_json_rx_browser(vhd, pss, in, len)) { + if (saiw_ws_json_rx_browser(vhd, pss, in, len, (lws_is_first_fragment(wsi) ? LWSSS_FLAG_SOM : 0) | + (lws_is_final_fragment(wsi) ? LWSSS_FLAG_EOM : 0))) { lwsl_wsi_err(wsi, "Closing because saiw_ws_json_rx_browser returned it"); return -1; diff --git a/src/web/w-private.h b/src/web/w-private.h index 6cd0817..f0da4cb 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -276,7 +276,7 @@ lws_struct_map_set(const lws_struct_map_t *map, char *u); int saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, - uint8_t *buf, size_t bl); + uint8_t *buf, size_t bl, unsigned int ss_flags); void sai_task_uuid_to_event_uuid(char *event_uuid33, const char *task_uuid65); @@ -300,7 +300,7 @@ int saiw_task_cancel(struct vhd *vhd, const char *task_uuid); int -saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len); +saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags); int saiw_get_blob(struct vhd *vhd, const char *url, sqlite3 **pdb, diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c index ead565f..9372dbf 100644 --- a/src/web/w-websrv.c +++ b/src/web/w-websrv.c @@ -34,7 +34,7 @@ typedef struct saiw_websrv { lws_struct_args_t a; struct lejp_ctx ctx; - struct lws_buflist *bltx; + struct lws_buflist *wbltx; } saiw_websrv_t; extern const lws_struct_map_t lsm_schema_json_map[]; @@ -50,6 +50,28 @@ enum { SAIS_WS_WEBSRV_RX_TASKACTIVITY, }; +int +saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) +{ + saiw_websrv_t *m = (saiw_websrv_t *)lws_ss_to_user_object(h); + unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); + + *pi = ss_flags; + + // lwsl_ss_notice(h, "sai-web: Queuing from browser -> sai-server"); + // lwsl_hexdump_notice(buf, len); + + if (lws_buflist_append_segment(&m->wbltx, buf - sizeof(int), len + sizeof(int)) < 0) { + lwsl_ss_err(h, "failed to append"); + return 1; + } + + if (lws_ss_request_tx(h)) + lwsl_ss_err(h, "failed to request tx"); + + return 0; +} + /* * sai-web is receiving from sai-server * @@ -203,32 +225,51 @@ 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, final = 1; - size_t fsl; - int used; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->wbltx), depi; + char som, som1, eom, final = 1; + size_t fsl, used; - if (!m->bltx) + if (!m->wbltx) { + lwsl_notice("%s: nothing to send from web -> srv\n", __func__); return LWSSSSRET_TX_DONT_SEND; + } + + depi = *pi; /* * We can only issue *len at a time. + * + * 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. */ - fsl = lws_buflist_next_segment_len(&m->bltx, NULL); + fsl = lws_buflist_next_segment_len(&m->wbltx, NULL); - if (fsl > *len) - final = 0; + lws_buflist_fragment_use(&m->wbltx, NULL, 0, &som, &eom); + if (som) { + fsl -= sizeof(int); + lws_buflist_fragment_use(&m->wbltx, buf, sizeof(int), &som1, &eom); + } - used = lws_buflist_fragment_use(&m->bltx, buf, *len, &som, &eom); + /* this is the only buflist user on pss->raw_tx */ + used = (size_t)lws_buflist_fragment_use(&m->wbltx, (uint8_t *)buf, *len, &som1, &eom); if (!used) return LWSSSSRET_TX_DONT_SEND; + if (used < fsl || (depi & LWS_WRITE_NO_FIN)) + final = 0; + + *len = used; *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); + lwsl_ss_notice(m->ss, "Sending %d ssflags %d", (int)*len, (int)*flags); - if (m->bltx) + if (m->wbltx) return lws_ss_request_tx(m->ss); return 0; @@ -253,7 +294,7 @@ saiw_lp_state(void *userobj, void *sh, lws_ss_constate_t state, return lws_ss_request_tx(m->ss); case LWSSSCS_DISCONNECTED: - lws_buflist_destroy_all_segments(&m->bltx); + lws_buflist_destroy_all_segments(&m->wbltx); lwsac_detach(&vhd->builders); break; @@ -280,20 +321,4 @@ const lws_ss_info_t ssi_saiw_websrv = { .streamtype = "websrv" }; -/* - * send to server - */ - -int -saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len) -{ - saiw_websrv_t *m = (saiw_websrv_t *)lws_ss_to_user_object(h); - - lwsl_ss_notice(h, "sai-web: Queuing from browser -> sai-server"); - lwsl_hexdump_notice(buf, len); - if (lws_buflist_append_segment(&m->bltx, buf, len) < 0) - return 1; - - return !!lws_ss_request_tx(h); -} diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 0abb878..ecd67f4 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -45,12 +45,16 @@ saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int struct pss *pss = lws_container_of(p, struct pss, same); int *pi = (int *)((const char *)buf - sizeof(int)); - if (pss->js_api_version >= api_ver_min) { + // if (pss->js_api_version >= api_ver_min) + { eff++; *pi = (int)flags; - if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) > 0) + if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) < 0) + lwsl_wsi_err(pss->wsi, "unable to buflist_append"); + else { lws_callback_on_writable(pss->wsi); + } } } lws_end_foreach_dll(p); @@ -437,7 +441,7 @@ saiw_event_state_change(struct vhd *vhd, const char *event_uuid) int saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, - size_t bl) + size_t bl, unsigned int ss_flags) { sai_browse_rx_taskinfo_t *ti; sai_browse_rx_evinfo_t *ei; @@ -529,7 +533,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, ei = (sai_browse_rx_evinfo_t *)a.dest; - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl); + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_EVENTRESET: @@ -546,7 +550,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsl_notice("%s: received request to reset event %s\n", __func__, ei->event_hash); - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl); + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_EVENTDELETE: @@ -562,7 +566,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsl_notice("%s: received request to delete event %s\n", __func__, ei->event_hash); - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl); + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); lwsac_free(&a.ac); break; @@ -592,7 +596,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, * User is asking us to rebuild a builder */ - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl); + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_PLATRESET: @@ -603,7 +607,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, * User is asking us to reset / rebuild a whole platform */ - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl); + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; default: @@ -1345,7 +1349,7 @@ send_it: (pss->send_state == WSS_SEND_ARTIFACT_INFO && (!sch || !sch->owner.head))) { /* does he want to subscribe to logs? */ - if (sch && sch->logsub && sch->one_task) { + if (sch && sch->logsub && sch->one_task && !pss->subs_list.owner) { strcpy(pss->sub_task_uuid, sch->one_task->uuid); lws_dll2_add_head(&pss->subs_list, &pss->vhd->subs_owner); pss->sub_timestamp = pss->initial_log_timestamp; /* where we got up to */ @@ -1430,5 +1434,5 @@ saiw_update_viewer_count(struct vhd *vhd) lws_struct_json_serialize_destroy(&js); if (len > 0) - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf + LWS_PRE, len); + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf + LWS_PRE, len, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); }