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 <string.h>
#include <signal.h>
#include <time.h>
+#include <assert.h>
#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);
}