diff --git a/assets/sai.css b/assets/sai.css
index 039b69e..7aeef9c 100644
--- a/assets/sai.css
+++ b/assets/sai.css
@@ -693,3 +693,15 @@ img.branch {
vertical-align:middle;
color: rgba(0, 0, 0, 0);
}
+
+.inst_box {
+ display: inline-block;
+ width: 12px;
+ height: 12px;
+ border: 1px solid #777;
+ margin-left: 2px;
+ vertical-align: middle;
+}
+.inst_idle { background-color: #d0d0d0; }
+.inst_run { background-color: #78c678; }
+.load_text { font-size: 9px; vertical-align: middle; }
diff --git a/assets/sai.js b/assets/sai.js
index fd02ec6..3331b95 100644
--- a/assets/sai.js
+++ b/assets/sai.js
@@ -826,6 +826,46 @@ function sai_event_render(o, now_ut, reset_all_icon)
return s;
}
+function getBuilderHostname(platName) {
+ return platName.split('.')[0];
+}
+
+function getBuilderGroupKey(platName) {
+ let hostname = platName.split('.')[0];
+ if (hostname.includes('-')) {
+ let parts = hostname.split('-');
+ return parts[parts.length - 1];
+ }
+ return hostname;
+}
+
+function createBuilderDiv(plat) {
+ const platDiv = document.createElement("div");
+ platDiv.className = "ibuil bdr";
+ platDiv.id = "binfo-" + plat.name;
+ platDiv.title = plat.platform + "@" + plat.name.split('.')[0] + " / " + plat.peer_ip;
+
+ let plat_parts = plat.platform.split('/');
+ let plat_os = plat_parts[0] || 'generic';
+ let plat_arch = plat_parts[1] || 'generic';
+ let plat_tc = plat_parts[2] || 'generic';
+
+ let innerHTML = `<table class="nomar"><tbody><tr><td class="bn">`;
+ innerHTML += `<img class="ip1 zup" src="/sai/${plat_os}.svg" onerror="this.src='/sai/generic.svg';this.onerror=null;">`;
+ innerHTML += `<img class="ip1 tread1" src="/sai/arch-${plat_arch}.svg" onerror="this.src='/sai/generic.svg';this.onerror=null;">`;
+ innerHTML += `<img class="ip1 tread2" src="/sai/tc-${plat_tc}.svg" onerror="this.src='/sai/generic.svg';this.onerror=null;">`;
+ innerHTML += `<br>${plat.peer_ip}`;
+ innerHTML += `<div id="instload-${plat.name}">`;
+
+ for (let i = 0; i < plat.instances; i++) {
+ innerHTML += `<div class="inst_box inst_idle" title="instance ${i}: idle"></div>`;
+ }
+
+ innerHTML += `</div></td></tr></tbody></table>`;
+ platDiv.innerHTML = innerHTML;
+ return platDiv;
+}
+
function render_builders(jso)
{
var s, n, conts = [], nc = 0;
@@ -877,13 +917,22 @@ function render_builders(jso)
if (!n ||
host !== jso.builders[n - 1].name.split('.')[0] ||
e.platform != samplat) {
- s += "<div class=\"ibuil bdr\" title=\"" +
+
+ s += "<div class=\"ibuil bdr\" title=\"" +
san(e.platform) + "@" + san(host) +
- (e.peer_ip ? " / " + san(e.peer_ip) : "") +
- "\"><table class=\"nomar\"><tr><td class=\"bn\">" +
+ (e.peer_ip ? " / " + san(e.peer_ip) : "") +
+ "\" id=\"binfo-" + san(e.name) + "\">" +
+ "<table class=\"nomar\"><tr><td class=\"bn\">" +
sai_plat_icon(e.platform, 1) +
(e.peer_ip ? "<br>" + san(e.peer_ip) : "");
+ /* Add a container for the instance load boxes */
+ s += "<div id=\"instload-" + san(e.name) + "\">";
+ for (var i = 0; i < e.instances; i++) {
+ s += "<div class=\"inst_box inst_idle\" title=\"instance " + i + ": idle\"></div>";
+ }
+ s += "</div>";
+
samplat = e.platform;
did = 1;
}
@@ -1064,15 +1113,17 @@ function ws_open_sai()
}
}
- if (jso.schema == "com.warmcat.sai.builders") {
+ switch (jso.schema) {
+
+ case "com.warmcat.sai.builders":
s = render_builders(jso);
if (document.getElementById("sai_builders"))
document.getElementById("sai_builders").innerHTML = s;
- }
+ break;
- if (jso.schema == "sai.warmcat.com.overview") {
+ case "sai.warmcat.com.overview":
/*
* Sent with an array of e[] to start, but also
* can send a single e[] if it just changed
@@ -1208,10 +1259,10 @@ function ws_open_sai()
sai.send(rs);
setTimeout(after_delete, 750);
});
- }
- }
+ }
+ break;
- if (jso.schema == "com.warmcat.sai.taskinfo") {
+ case "com.warmcat.sai.taskinfo":
authd = jso.authorized;
if (jso.authorized === 0) {
@@ -1348,10 +1399,116 @@ function ws_open_sai()
aging();
}
+ break;
+
+
+ case "sai-builders":
+ const buildersContainer = document.getElementById("sai_builders");
+ if (!buildersContainer) { break; }
+ buildersContainer.innerHTML = "";
+ if (!jso.platforms || !Array.isArray(jso.platforms) || jso.platforms.length === 0) {
+ break;
+ }
+
+ // --- NEW, FINAL LOGIC ---
+
+ // Step 1: Create a map of group keys to group container divs.
+ // Also create a list of standalone platforms.
+ const groups = {};
+ const standalones = [];
+ const platformsByGroup = {};
+
+ for (const plat of jso.platforms) {
+ const groupKey = getBuilderGroupKey(plat.name);
+ if (!platformsByGroup[groupKey]) {
+ platformsByGroup[groupKey] = [];
}
+ platformsByGroup[groupKey].push(plat);
+ }
+
+ // Step 2: Build the HTML structure
+ const table = document.createElement("table");
+ table.className = "builders";
+ const tbody = document.createElement("tbody");
+ const tr = document.createElement("tr");
+ const td = document.createElement("td");
+ tr.appendChild(td);
+ tbody.appendChild(tr);
+ table.appendChild(tbody);
+
+ // Order of rendering matters. Let's process the groups first.
+ const groupKeys = Object.keys(platformsByGroup).sort();
+
+ for (const key of groupKeys) {
+ const groupPlatforms = platformsByGroup[key];
-
- if (jso.schema == "com-warmcat-sai-artifact") {
+ // Find platforms that are nested (VMs, different hostnames)
+ const nestedPlatforms = groupPlatforms.filter(p => getBuilderHostname(p.name) !== key);
+ // Find platforms that are the main host itself
+ const mainPlatforms = groupPlatforms.filter(p => getBuilderHostname(p.name) === key);
+
+ // If there are nested platforms, create a group container for them.
+ if (nestedPlatforms.length > 0) {
+ const groupDiv = document.createElement("div");
+ groupDiv.className = "ibuil ibuilctr bdr";
+
+ const nameDiv = document.createElement("div");
+ nameDiv.className = "ibuilctrname bdr";
+ nameDiv.textContent = key;
+ groupDiv.appendChild(nameDiv);
+
+ for (const plat of nestedPlatforms) {
+ groupDiv.appendChild(createBuilderDiv(plat));
+ }
+ td.appendChild(groupDiv);
+ }
+
+ // Render the main host platforms as standalones.
+ for (const plat of mainPlatforms) {
+ td.appendChild(createBuilderDiv(plat));
+ }
+ }
+
+ buildersContainer.appendChild(table);
+ break;
+
+ case "com.warmcat.sai.loadreport":
+ if (!jso.platform_name)
+ break;
+
+ const platformName = jso.platform_name;
+
+ // The container for the load squares has a predictable ID
+ const loadContainerId = "instload-" + platformName;
+ const loadContainer = document.getElementById(loadContainerId);
+
+ if (loadContainer) {
+
+ // Clear any old load indicators
+ loadContainer.innerHTML = "";
+
+ // Loop through the new load data and create the squares
+ if (jso.loads && Array.isArray(jso.loads)) {
+
+ for (const instanceLoad of jso.loads) {
+ let instanceDiv = document.createElement("div");
+
+ // Set class for styling (e.g., green for idle, red for busy)
+ let stateClass = instanceLoad.state ? "inst_busy" : "inst_idle";
+ instanceDiv.className = "inst_box " + stateClass;
+
+ // Set tooltip to show the CPU percentage
+ let cpu = instanceLoad.cpu_percent / 10.0;
+ instanceDiv.title = `Instance: ${stateClass.split('_')[1]} \nCPU: ${cpu.toFixed(1)}%`;
+
+ loadContainer.appendChild(instanceDiv);
+ }
+ }
+ }
+
+ break;
+
+ case "com-warmcat-sai-artifact":
console.log(jso);
sai_arts += "<div class=\"sai_arts\"><img src=\"artifact.svg\"> <a href=\"artifacts/" +
@@ -1364,9 +1521,9 @@ function ws_open_sai()
if (document.getElementById("sai_arts"))
document.getElementById("sai_arts").innerHTML = sai_arts;
- }
+ break;
- if (jso.schema == "com-warmcat-sai-logs") {
+ case "com-warmcat-sai-logs":
var s1 = atob(jso.log), s = hsanitize(s1), li,
en = "", yo, dh, ce, tn = "";
@@ -1398,9 +1555,10 @@ function ws_open_sai()
}
- if (!cont[jso.channel] && jso.len)
+ if (cont && !cont[jso.channel] && jso.len)
tn = ((jso.timestamp - tfirst) / 1000000).toFixed(4);
+ if (cont)
cont[jso.channel] = (li == 0);
while (li--) {
@@ -1438,9 +1596,10 @@ function ws_open_sai()
document.body.clientHeight;
}, 500);
}
- }
- };
-
+
+ break;
+ } /* switch */
+ } /* onmessage */
sai.onclose = function(){
// document.getElementById("title").innerHTML =
// "Server Status (Disconnected)";
diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c
index 95cce9a..9e7cf97 100644
--- a/src/builder/b-comms.c
+++ b/src/builder/b-comms.c
@@ -27,6 +27,14 @@
#include "../common/struct-metadata.c"
+const lws_struct_map_t lsm_viewerstate_members[] = {
+ LSM_UNSIGNED(sai_viewer_state_t, viewers, "viewers"),
+};
+
+static const lws_struct_map_t lsm_schema_json_loadreport[] = {
+ LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, "com.warmcat.sai.loadreport"),
+};
+
static lws_ss_state_return_t
saib_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
{
@@ -257,6 +265,44 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
}
/*
+ * Any load reports to send?
+ */
+ if (spm->load_report_owner.count) {
+ struct lws_dll2 *d = lws_dll2_get_head(&spm->load_report_owner);
+ sai_load_report_t *lr =
+ lws_container_of(d, sai_load_report_t, list);
+
+ // lwsl_notice("%s: issuing load report for %s\n", __func__,
+ // lr->builder_name);
+
+ js = lws_struct_json_serialize_create(lsm_schema_json_loadreport,
+ LWS_ARRAY_SIZE(lsm_schema_json_loadreport), 0, lr);
+ if (!js)
+ return -1;
+
+ n = (int)lws_struct_json_serialize(js, start,
+ lws_ptr_diff_size_t(end, start), &w);
+ lws_struct_json_serialize_destroy(&js);
+
+ // lwsl_hexdump_notice(start, w);
+
+ n = (int)w;
+
+ lws_dll2_remove(&lr->list);
+ lws_start_foreach_dll_safe(struct lws_dll2 *, il, il1, lr->loads.head) {
+ sai_instance_load_t *i = lws_container_of(il, sai_instance_load_t, list);
+ lws_dll2_remove(&i->list);
+ free(i);
+ } lws_end_foreach_dll_safe(il, il1);
+ free(lr);
+
+ r = lws_ss_request_tx(spm->ss);
+ if (r)
+ return r;
+ goto sendify;
+ }
+
+ /*
* Any resource requests / relinquishments to process?
*/
@@ -536,6 +582,60 @@ cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user)
return 0;
}
+void
+saib_sul_load_report_cb(struct lws_sorted_usec_list *sul)
+{
+ struct sai_plat_server *spm = lws_container_of(sul,
+ struct sai_plat_server, sul_load_report);
+ sai_load_report_t *lr = calloc(1, sizeof(*lr));
+ struct sai_plat *sp = NULL;
+
+ if (!lr)
+ return;
+
+ /*
+ * This builder process has one name, but may have multiple platforms,
+ * each with multiple instances. For now, we report on the whole builder
+ * under one name.
+ */
+ lws_strncpy(lr->builder_name, builder.host, sizeof(lr->builder_name));
+ if (builder.sai_plat_owner.head) {
+ sai_plat_t *any_plat = lws_container_of(builder.sai_plat_owner.head,
+ sai_plat_t, sai_plat_list);
+ lws_strncpy(lr->platform_name, any_plat->name, sizeof(lr->platform_name));
+ }
+
+ /*
+ * Iterate all platforms and their nspawn instances to collect load.
+ * In a real implementation, you would query the system for CPU usage
+ * of each nspawn process. For now, we will simulate it.
+ */
+ lws_start_foreach_dll(struct lws_dll2 *, p, builder.sai_plat_owner.head) {
+ sp = lws_container_of(p, sai_plat_t, sai_plat_list);
+
+ lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) {
+ struct sai_nspawn *ns = lws_container_of(d, struct sai_nspawn, list);
+ sai_instance_load_t *il = calloc(1, sizeof(*il));
+
+ if (il) {
+ il->state = (ns->state == NSSTATE_BUILD);
+ /* Simulate load: 50% if building, 1% if idle */
+ il->cpu_percent = il->state ? 500 : 10;
+ lws_dll2_add_tail(&il->list, &lr->loads);
+ }
+ } lws_end_foreach_dll(d);
+ } lws_end_foreach_dll(p);
+
+ lws_dll2_add_tail(&lr->list, &spm->load_report_owner);
+ if (lws_ss_request_tx(spm->ss))
+ lwsl_debug("%s: request tx failed\n", __func__);
+
+ /* Reschedule the timer */
+ lws_sul_schedule(builder.context, 0, &spm->sul_load_report,
+ saib_sul_load_report_cb, SAI_LOAD_REPORT_US);
+}
+
+
static lws_ss_state_return_t
saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
lws_ss_tx_ordinal_t ack)
@@ -622,6 +722,9 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
case LWSSSCS_CONNECTED:
lwsl_user("%s: CONNECTED: %p\n", __func__, spm->ss);
spm->phase = PHASE_START_ATTACH;
+ /* Initialize the load report SUL timer for this server connection */
+ lws_sul_cancel(&spm->sul_load_report);
+
return lws_ss_request_tx(spm->ss);
case LWSSSCS_DISCONNECTED:
@@ -630,6 +733,7 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state,
*/
lwsl_user("%s: DISCONNECTED\n", __func__);
+ lws_sul_cancel(&spm->sul_load_report);
lws_dll2_foreach_safe(&builder.sai_plat_owner, spm,
cleanup_on_ss_disconnect);
break;
diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c
index b595aa0..fca417e 100644
--- a/src/builder/b-nspawn.c
+++ b/src/builder/b-nspawn.c
@@ -30,6 +30,8 @@
#include <sys/types.h>
#if !defined(WIN32)
#include <unistd.h>
+#else
+#include <process.h>
#endif
#include <fcntl.h>
#include <assert.h>
@@ -243,9 +245,8 @@ int
saib_spawn(struct sai_nspawn *ns)
{
struct lws_spawn_piped_info info;
- char args[290], st[2048], cgroup[128], *p;
+ char args[290], st[2048], *p;
const char *respath = "unk";
- int fd, n, in_cgroup = 1;
const char * cmd[] = {
"/bin/ps",
NULL
@@ -255,6 +256,11 @@ saib_spawn(struct sai_nspawn *ns)
"LANG=en_US.UTF-8",
NULL
};
+ int fd, n;
+#if defined(__linux__)
+ int in_cgroup = 1;
+ char cgroup[128];
+#endif
lws_strncpy(st, ns->sp->name, sizeof(st));
lws_filename_purify_inplace(st);
@@ -318,7 +324,9 @@ saib_spawn(struct sai_nspawn *ns)
cmd[0] = args;
+#if defined(__linux__)
lws_snprintf(cgroup, sizeof(cgroup), "inst-%u-%d", (unsigned int)getpid(), ns->instance_idx);
+#endif
memset(&info, 0, sizeof(info));
info.vh = builder.vhost;
@@ -330,8 +338,10 @@ saib_spawn(struct sai_nspawn *ns)
info.reap_cb = sai_lsp_reap_cb;
info.opaque = ns;
info.plsp = &ns->lsp;
+#if defined(__linux__)
info.cgroup_name_suffix = cgroup;
info.p_cgroup_ret = &in_cgroup;
+#endif
ns->lsp = lws_spawn_piped(&info);
if (!ns->lsp) {
@@ -340,7 +350,9 @@ saib_spawn(struct sai_nspawn *ns)
return 1;
}
+#if defined(__linux__)
lwsl_notice("%s: lws_spawn_piped started (cgroup: %d)\n", __func__, in_cgroup);
+#endif
return 0;
}
diff --git a/src/builder/b-private.h b/src/builder/b-private.h
index 38e2bae..e37159f 100644
--- a/src/builder/b-private.h
+++ b/src/builder/b-private.h
@@ -39,6 +39,7 @@
#include <pthread.h>
#include <git2.h>
+#define SAI_LOAD_REPORT_US (5 * LWS_US_PER_SEC)
#define SAI_IDLE_GRACE_US (30 * LWS_US_PER_SEC)
#define SAI_STAY_POLL_US (20 * LWS_US_PER_SEC)
@@ -221,3 +222,6 @@ saib_create_resproxy_listen_uds(struct lws_context *context,
int
saib_handle_resource_result(struct sai_plat_server *spm, const char *in, size_t len);
+
+void
+saib_sul_load_report_cb(struct lws_sorted_usec_list *sul);
diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c
index 26e34d0..d361a1d 100644
--- a/src/builder/b-sai.c
+++ b/src/builder/b-sai.c
@@ -32,8 +32,10 @@
#include <stdlib.h>
#include <sys/types.h>
+#if !defined(WIN32)
#include <pwd.h>
#include <grp.h>
+#endif
#if defined(__linux__)
#include <unistd.h>
diff --git a/src/builder/b-task.c b/src/builder/b-task.c
index 14257bc..c827217 100644
--- a/src/builder/b-task.c
+++ b/src/builder/b-task.c
@@ -33,16 +33,23 @@ static char csep = '/';
static char csep = '\\';
#endif
+static const lws_struct_map_t lsm_viewerstate_members[] = {
+ LSM_UNSIGNED(sai_viewer_state_t, viewers, "viewers"),
+};
+
const lws_struct_map_t lsm_schema_map_m_to_b[] = {
- LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"),
- LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, "com.warmcat.sai.taskcan"),
- LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, "com-warmcat-sai-resource")
+ LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"),
+ LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, "com.warmcat.sai.taskcan"),
+ LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewerstate_members,
+ "com.warmcat.sai.viewerstate"),
+ LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, "com-warmcat-sai-resource")
};
enum {
SAIB_RX_TASK_ALLOCATION,
SAIB_RX_TASK_CANCEL,
- SAIB_RX_RESOURCE_REPLY,
+ SAIB_RX_VIEWERSTATE,
+ SAIB_RX_RESOURCE_REPLY
};
static const char * const nsstates[] = {
@@ -770,6 +777,23 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len)
} lws_end_foreach_dll_safe(mp, mp1);
break;
+ case SAIB_RX_VIEWERSTATE:
+ {
+ sai_viewer_state_t *vs = (sai_viewer_state_t *)a.dest;
+ lwsl_notice("Received viewer state update: %u viewers\n",
+ vs->viewers);
+
+ spm->viewer_count = vs->viewers;
+
+ if (vs->viewers)
+ /* At least one viewer, start reporting */
+ lws_sul_schedule(builder.context, 0, &spm->sul_load_report,
+ saib_sul_load_report_cb, SAI_LOAD_REPORT_US);
+ else
+ lws_sul_cancel(&spm->sul_load_report);
+ }
+ break;
+
case SAIB_RX_RESOURCE_REPLY:
reso = (sai_resource_t *)a.dest;
diff --git a/src/common/include/private.h b/src/common/include/private.h
index 5378226..5291eed 100644
--- a/src/common/include/private.h
+++ b/src/common/include/private.h
@@ -54,6 +54,35 @@ enum {
SAISPRF_SIGNALLED = 0x4000,
};
+/*
+ * per-instance load data.
+ * Sent from builder -> server -> web -> browser.
+ */
+typedef struct sai_instance_load {
+ lws_dll2_t list;
+ unsigned int cpu_percent; /* CPU usage for this instance * 100 */
+ unsigned int state; /* 0 = idle, 1 = running */
+} sai_instance_load_t;
+
+/*
+ * load report struct, sent in its own schema.
+ */
+typedef struct sai_load_report {
+ lws_dll2_t list; /* Not used, for schema mapping */
+ char builder_name[64];
+ char platform_name[128];
+ lws_dll2_owner_t loads;
+} sai_load_report_t;
+
+/*
+ * viewer state.
+ * Sent from server -> builder.
+ */
+typedef struct sai_viewer_state {
+ lws_dll2_t list; /* Not used, for schema mapping */
+ unsigned int viewers;
+} sai_viewer_state_t;
+
struct sai_nspawn;
@@ -297,6 +326,11 @@ typedef struct sai_plat_server {
char resproxy_path[128];
+ /* for load reporting */
+ lws_sorted_usec_list_t sul_load_report;
+ lws_dll2_owner_t load_report_owner; /* sai_load_report_t */
+ unsigned int viewer_count;
+
const char *url;
const char *name;
@@ -356,6 +390,7 @@ typedef struct sai_plat {
struct lws *wsi; /* server side only */
lws_dll2_owner_t env_head;
+ lws_dll2_owner_t loads;
int instances;
int ongoing;
diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c
index cff297d..560e6ac 100644
--- a/src/common/struct-metadata.c
+++ b/src/common/struct-metadata.c
@@ -21,6 +21,19 @@
* lws_struct metadata for structs common to builder and server
*/
+const lws_struct_map_t lsm_instance_load[] = {
+ LSM_UNSIGNED (sai_instance_load_t, cpu_percent, "cpu_percent"),
+ LSM_UNSIGNED (sai_instance_load_t, state, "state"),
+};
+
+const lws_struct_map_t lsm_load_report_members[] = {
+ LSM_CARRAY (sai_load_report_t, builder_name, "builder_name"),
+ LSM_CARRAY (sai_load_report_t, platform_name, "platform_name"),
+ LSM_LIST (sai_load_report_t, loads, sai_instance_load_t, list,
+ NULL, lsm_instance_load, "loads"),
+};
+
+
static const lws_struct_map_t lsm_plat[] = {
LSM_STRING_PTR (sai_plat_t, name, "name"),
LSM_UNSIGNED (sai_plat_t, ongoing, "ongoing"),
diff --git a/src/power/p-intake.c b/src/power/p-intake.c
index f59a52e..bb2d386 100644
--- a/src/power/p-intake.c
+++ b/src/power/p-intake.c
@@ -144,6 +144,8 @@ callback_ws_power(struct lws *wsi, enum lws_callback_reasons reason, void *user,
* Update the sai-webs about the builder removal, so they
* can update their connected browsers
*/
+ lwsl_wsi_warn(wsi, "LWS_CALLBACK_CLOSED: doing WSS_PREPARE_BUILDER_SUMMARY\n");
+
sais_list_builders(vhd);
break;
diff --git a/src/server/s-comms.c b/src/server/s-comms.c
index c72ce27..9fe5685 100644
--- a/src/server/s-comms.c
+++ b/src/server/s-comms.c
@@ -776,6 +776,7 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user,
* Update the sai-webs about the builder removal, so they
* can update their connected browsers
*/
+ lwsl_wsi_warn(pss->wsi, "LWS_CALLBACK_CLOSED: doing WSS_PREPARE_BUILDER_SUMMARY\n");
sais_list_builders(vhd);
break;
@@ -791,12 +792,14 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user,
pss->wsi = wsi;
if (sais_ws_json_rx_builder(vhd, pss, in, len))
return -1;
+
if (!pss->announced) {
/*
- * Update the sai-webs about the builder removal, so
+ * Update the sai-webs about the builder creation, so
* they can update their connected browsers
*/
+ lwsl_wsi_warn(pss->wsi, "LWS_CALLBACK_RECEIVE: unannounced pss doing WSS_PREPARE_BUILDER_SUMMARY\n");
sais_list_builders(vhd);
pss->announced = 1;
diff --git a/src/server/s-private.h b/src/server/s-private.h
index 1afe772..c985720 100644
--- a/src/server/s-private.h
+++ b/src/server/s-private.h
@@ -106,6 +106,7 @@ struct pss {
lws_dll2_owner_t res_pending_reply_owner; /* sai_resource_msg_t
* resource JSON return
* messages to builder */
+ lws_dll2_owner_t viewer_state_owner;
lws_struct_args_t a;
union {
@@ -198,6 +199,8 @@ struct vhd {
const char *notification_key;
+ unsigned int browser_viewer_count;
+
sais_t server;
};
diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c
index ec0735a..ffb3afb 100644
--- a/src/server/s-websrv.c
+++ b/src/server/s-websrv.c
@@ -47,13 +47,17 @@ typedef struct websrvss_srv {
struct lejp_ctx ctx;
struct lws_buflist *bltx;
-
+ unsigned int viewers;
} websrvss_srv_t;
static lws_struct_map_t lsm_browser_taskreset[] = {
LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "uuid"),
};
+static const lws_struct_map_t lsm_viewercount_members[] = {
+ LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"),
+};
+
static const lws_struct_map_t lsm_schema_json_map[] = {
LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset,
/* shares struct */ "com.warmcat.sai.taskreset"),
@@ -63,13 +67,16 @@ static const lws_struct_map_t lsm_schema_json_map[] = {
/* shares struct */ "com.warmcat.sai.eventdelete"),
LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel,
"com.warmcat.sai.taskcan"),
+ LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewercount_members,
+ "com.warmcat.sai.viewercount"),
};
enum {
SAIS_WS_WEBSRV_RX_TASKRESET,
SAIS_WS_WEBSRV_RX_EVENTRESET,
SAIS_WS_WEBSRV_RX_EVENTDELETE,
- SAIS_WS_WEBSRV_RX_TASKCANCEL
+ SAIS_WS_WEBSRV_RX_TASKCANCEL,
+ SAIS_WS_WEBSRV_RX_VIEWERCOUNT,
};
int
@@ -275,6 +282,12 @@ sais_eventchange(struct lws_ss_handle *hsrv, const char *event_uuid, int state)
lws_ss_server_foreach_client(hsrv, _sais_eventchange, (void *)&arg);
}
+static void
+sum_viewers_cb(struct lws_ss_handle *h, void *arg)
+{
+ websrvss_srv_t *m_client = (websrvss_srv_t *)lws_ss_to_user_object(h);
+ *(unsigned int *)arg += m_client->viewers;
+}
static lws_ss_state_return_t
websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
@@ -467,6 +480,39 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
sais_task_cancel(m->vhd, ei->event_hash);
break;
+
+ case SAIS_WS_WEBSRV_RX_VIEWERCOUNT:
+ {
+ sai_viewer_state_t *vs = (sai_viewer_state_t *)a.dest;
+ unsigned int total_viewers = 0;
+
+ /* Store viewer count for this specific sai-web client */
+ m->viewers = vs->viewers;
+
+ /* Recalculate total from all connected sai-web clients */
+ lws_ss_server_foreach_client(m->vhd->h_ss_websrv,
+ sum_viewers_cb, &total_viewers);
+
+ m->vhd->browser_viewer_count = total_viewers;
+ lwsl_notice("%s: Client viewer count %u, total is now %u\n",
+ __func__, m->viewers, m->vhd->browser_viewer_count);
+
+ /* Broadcast the new viewer state to all connected builders */
+ lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) {
+ struct pss *pss_builder = lws_container_of(p,
+ struct pss, same);
+ sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend));
+
+ if (vsend) {
+ /* Send the new TOTAL viewer count */
+ vsend->viewers = m->vhd->browser_viewer_count;
+ lws_dll2_add_tail(&vsend->list,
+ &pss_builder->viewer_state_owner);
+ lws_callback_on_writable(pss_builder->wsi);
+ }
+ } lws_end_foreach_dll(p);
+ break;
+ }
}
return 0;
@@ -511,13 +557,43 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state,
// lws_ss_state_name((int)state), (unsigned int)ack);
switch (state) {
- case LWSSSCS_DISCONNECTED:
- lws_buflist_destroy_all_segments(&m->bltx);
+ case LWSSSCS_DISCONNECTED: {
+ unsigned int total_viewers = 0;
+
+ lws_buflist_destroy_all_segments(&m->bltx);
+ m->viewers = 0;
+
+ /* This sai-web client disconnected, recalculate total viewers */
+ lws_ss_server_foreach_client(m->vhd->h_ss_websrv,
+ sum_viewers_cb, &total_viewers);
+
+ if (m->vhd->browser_viewer_count != total_viewers) {
+ m->vhd->browser_viewer_count = total_viewers;
+ lwsl_notice("%s: A sai-web client disconnected, total viewers now %u\n",
+ __func__, total_viewers);
+
+ /* Broadcast new count to builders */
+ lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) {
+ struct pss *pss_builder = lws_container_of(p,
+ struct pss, same);
+ sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend));
+
+ if (vsend) {
+ vsend->viewers = m->vhd->browser_viewer_count;
+ lws_dll2_add_tail(&vsend->list,
+ &pss_builder->viewer_state_owner);
+ lws_callback_on_writable(pss_builder->wsi);
+ }
+ } lws_end_foreach_dll(p);
+ }
break;
+ }
case LWSSSCS_CREATING:
+ m->viewers = 0;
return lws_ss_request_tx(m->ss);
case LWSSSCS_CONNECTED:
+ // lwsl_warn("%s: resending builders because CONNECTED\n", __func__);
sais_list_builders(m->vhd);
break;
case LWSSSCS_ALL_RETRIES_FAILED:
diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c
index 15ae7ea..17bc2f7 100644
--- a/src/server/s-ws-builder.c
+++ b/src/server/s-ws-builder.c
@@ -40,6 +40,20 @@ typedef struct sais_logcache_pertask {
lws_dll2_owner_t cache; /* sai_log_t */
} sais_logcache_pertask_t;
+/* map for a single instance's load */
+static const lws_struct_map_t lsm_instance_load[] = {
+ LSM_UNSIGNED(sai_instance_load_t, cpu_percent, "cpu_percent"),
+ LSM_UNSIGNED(sai_instance_load_t, state, "state"),
+};
+
+/* map for the members of the load report object */
+static const lws_struct_map_t lsm_load_report_members[] = {
+ LSM_CARRAY(sai_load_report_t, builder_name, "builder_name"),
+ LSM_CARRAY(sai_load_report_t, platform_name, "platform_name"),
+ LSM_LIST(sai_load_report_t, loads, sai_instance_load_t, list,
+ NULL, lsm_instance_load, "loads"),
+};
+
/*
* The Schema that may be sent to us by a builder
*
@@ -57,6 +71,8 @@ static const lws_struct_map_t lsm_schema_map_ba[] = {
"com.warmcat.sai.taskrej"),
LSM_SCHEMA (sai_artifact_t, NULL, lsm_artifact,
"com-warmcat-sai-artifact"),
+ LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, /* from builder */
+ "com.warmcat.sai.loadreport"),
LSM_SCHEMA (sai_resource_t, NULL, lsm_resource,
"com-warmcat-sai-resource"),
};
@@ -66,7 +82,8 @@ enum {
SAIM_WSSCH_BUILDER_LOGS,
SAIM_WSSCH_BUILDER_TASKREJ,
SAIM_WSSCH_BUILDER_ARTIFACT,
- SAIM_WSSCH_BUILDER_RESOURCE_REQ
+ SAIM_WSSCH_BUILDER_LOADREPORT,
+ SAIM_WSSCH_BUILDER_RESOURCE_REQ,
};
static void
@@ -368,6 +385,7 @@ handle:
cb->instances = build->instances;
cb->wsi = pss->wsi;
+ pss->announced = 0;
/* Then attach the copy to the server in the vhd
*/
@@ -486,7 +504,7 @@ bail:
cb->ongoing = rej->ongoing;
cb->instances = rej->limit;
- lwsl_notice("%s: builder %s reports load %d/%d (rej %s)\n",
+ lwsl_notice("%s: builder %s reports occupancy %d/%d (rej %s)\n",
__func__, cb->name, cb->ongoing, cb->instances,
rej->task_uuid[0] ? rej->task_uuid : "none");
@@ -496,6 +514,11 @@ bail:
lwsac_free(&pss->a.ac);
break;
+ case SAIM_WSSCH_BUILDER_LOADREPORT:
+ lwsl_wsi_user(pss->wsi, "SAIM_WSSCH_BUILDER_LOADREPORT broadcasting\n");
+ sais_websrv_broadcast(vhd->h_ss_websrv, (const char *)buf, bl);
+ break;
+
case SAIM_WSSCH_BUILDER_ARTIFACT:
/*
* Builder wants to send us an artifact.
@@ -814,6 +837,50 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf,
sai_task_t *task;
size_t w;
+ if (pss->viewer_state_owner.head) {
+ /*
+ * Pending viewer state message to send to a builder
+ */
+ sai_viewer_state_t *vs = lws_container_of(
+ pss->viewer_state_owner.head,
+ sai_viewer_state_t, list);
+
+ const lws_struct_map_t lsm_viewerstate_members[] = {
+ LSM_UNSIGNED(sai_viewer_state_t, viewers, "viewers"),
+ };
+ const lws_struct_map_t lsm_schema_viewerstate[] = {
+ LSM_SCHEMA(sai_viewer_state_t, NULL, lsm_viewerstate_members,
+ "com.warmcat.sai.viewerstate")
+ };
+
+ lwsl_wsi_notice(pss->wsi, "++++ Sending viewerstate (count: %u) to builder\n",
+ vs->viewers);
+
+ js = lws_struct_json_serialize_create(lsm_schema_viewerstate,
+ LWS_ARRAY_SIZE(lsm_schema_viewerstate), 0, vs);
+ if (!js)
+ return 1;
+
+ n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w);
+ lws_struct_json_serialize_destroy(&js);
+
+ /* Dequeue the message we just sent */
+ lws_dll2_remove(&vs->list);
+ /* And free the memory */
+ free(vs);
+
+ first = 1;
+
+ /*
+ * If there are more viewer state messages, or other messages,
+ * * request another writeable callback.
+ */
+ if (pss->viewer_state_owner.head)
+ lws_callback_on_writable(pss->wsi);
+
+ goto send_json;
+ }
+
if (pss->task_cancel_owner.head) {
/*
* Pending cancel message to send
@@ -832,9 +899,6 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf,
lws_dll2_remove(&c->list);
free(c);
- first = 1;
- pss->walk = NULL;
-
goto send_json;
}
@@ -858,9 +922,6 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf,
lws_dll2_remove(&rm->list);
free(rm);
- first = 1;
- pss->walk = NULL;
-
goto send_json;
}
@@ -891,9 +952,6 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf,
lwsac_free(&task->ac_task_container);
first = 1;
- pss->walk = NULL;
-
- //lwsac_free(&pss->query_ac);
send_json:
p += w;
@@ -915,7 +973,10 @@ send_json:
(enum lws_write_protocol)flags) < 0)
return -1;
- lws_callback_on_writable(pss->wsi);
+ if (pss->viewer_state_owner.head || pss->task_cancel_owner.head ||
+ pss->res_pending_reply_owner.count ||
+ pss->issue_task_owner.count)
+ lws_callback_on_writable(pss->wsi);
return 0;
}
diff --git a/src/web/w-comms.c b/src/web/w-comms.c
index 29880d0..c61b181 100644
--- a/src/web/w-comms.c
+++ b/src/web/w-comms.c
@@ -39,6 +39,8 @@
#include "../common/struct-metadata.c"
+extern const lws_struct_map_t lsm_schema_json_map[];
+
typedef enum {
SJS_CLONING,
SJS_ASSIGNING,
@@ -81,8 +83,45 @@ const lws_struct_map_t lsm_schema_sq3_map_auth[] = {
LSM_SCHEMA_DLL2 (sai_auth_t, list, NULL, lsm_auth, "auth"),
};
+static lws_struct_map_t lsm_websrv_evinfo[] = {
+ LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"),
+};
+
+const lws_struct_map_t lsm_schema_json_map[] = {
+ LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
+ /* shares struct */ "sai-taskchange"),
+ LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
+ /* shares struct */ "sai-eventchange"),
+ LSM_SCHEMA (sai_plat_owner_t, NULL, lsm_plat_list, "sai-builders"),
+ LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
+ /* shares struct */ "sai-overview"),
+ LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
+ /* shares struct */ "sai-tasklogs"),
+ LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members,
+ "com.warmcat.sai.loadreport"),
+};
+
+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
@@ -974,35 +1013,24 @@ clean_spa:
/*
* ws connections from builders and browsers
*/
-
- case LWS_CALLBACK_FILTER_PROTOCOL_CONNECTION:
- n = lws_hdr_copy(wsi, (char *)buf, sizeof(buf) - 1,
- WSI_TOKEN_GET_URI);
- if (!n)
- buf[0] = '\0';
- //lwsl_notice("%s: checking with lwsgs for ws conn: %s\n",
- // __func__, (const char *)buf);
-
- /*
- * Builders don't authenticate using sessions...
- */
-
- if (n >= 8 && !strncmp((const char *)buf + n - 8,
- "/builder", 8))
- return 0;
-
- return 0;
-#if 0
- /* but everything else does */
-
- car_args.max_len = LWSGS_AUTH_LOGGED_IN | LWSGS_AUTH_VERIFIED;
- car_args.final = 0;
- car_args.chunked = 1; /* ie, we are ws */
-
- in = &car_args;
- goto passthru;
-#endif
-
+ case LWS_CALLBACK_FILTER_PROTOCOL_CONNECTION:
+ n = lws_hdr_copy(wsi, (char *)buf, sizeof(buf) - 1,
+ WSI_TOKEN_GET_URI);
+
+ /*
+ * This protocol is for browsers on /browse... URLs.
+ * Builders connect on /builder... URLs and should be handled
+ * by a different protocol. Explicitly reject them here.
+ *
+ * Returning 0 accepts the connection for this protocol.
+ * Returning non-zero rejects it.
+ */
+ if (n >= 8 && !strncmp((const char *)buf + n - 8,
+ "/builder", 8))
+ return 1; /* Reject builder connections */
+
+ return 0;
+
case LWS_CALLBACK_ESTABLISHED:
if (!vhd) {
@@ -1023,8 +1051,6 @@ clean_spa:
ck.aud = vhd->jwt_audience;
ck.cookie_name = "__Host-sai_jwt";
- lws_dll2_add_head(&pss->same, &vhd->browsers);
-
cml = sizeof(buf);
if (!lws_jwt_get_http_cookie_validate_jwt(wsi, &ck,
(char *)buf, &cml) &&
@@ -1109,6 +1135,8 @@ clean_spa:
sizeof(pss->specific_ref));
}
+ saiw_browser_state_changed(pss, 1);
+
lwsl_info("%s: spec %d, ref '%s', task '%s' \n", __func__,
pss->specificity, pss->specific_ref, pss->specific_task);
break;
@@ -1116,6 +1144,7 @@ clean_spa:
if (!strcmp((char *)start, "/browse")) {
lwsl_info("%s: ESTABLISHED: browser\n", __func__);
+ saiw_browser_state_changed(pss, 1);
pss->wsi = wsi;
break;
}
@@ -1127,7 +1156,8 @@ clean_spa:
case LWS_CALLBACK_CLOSED:
lwsl_err("%s: CLOSED browse conn\n", __func__);
- lws_dll2_remove(&pss->same);
+ lws_buflist_destroy_all_segments(&pss->raw_tx);
+ saiw_browser_state_changed(pss, 0);
lws_dll2_remove(&pss->subs_list);
lws_dll2_foreach_safe(&pss->sched, NULL, saiw_sched_destroy);
diff --git a/src/web/w-private.h b/src/web/w-private.h
index dbc2eff..9fbef99 100644
--- a/src/web/w-private.h
+++ b/src/web/w-private.h
@@ -131,6 +131,7 @@ struct pss {
struct lws_spa *spa;
struct lejp_ctx ctx;
+ struct lws_buflist *raw_tx;
sai_notification_t sn;
struct lws_dll2 same; /* owner: vhd.browsers */
@@ -315,3 +316,13 @@ saiw_dealloc_sched(saiw_scheduled_t *sch);
int
saiw_sched_destroy(struct lws_dll2 *d, void *user);
+
+void
+saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len);
+
+void
+saiw_browser_state_changed(struct pss *pss, int established);
+
+void
+saiw_update_viewer_count(struct vhd *vhd);
+
diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c
index 1344e60..ced0742 100644
--- a/src/web/w-websrv.c
+++ b/src/web/w-websrv.c
@@ -39,30 +39,24 @@ typedef struct saiw_websrv {
} saiw_websrv_t;
-static lws_struct_map_t lsm_websrv_evinfo[] = {
- LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"),
-};
-
-static const lws_struct_map_t lsm_schema_json_map[] = {
- LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
- /* shares struct */ "sai-taskchange"),
- LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
- /* shares struct */ "sai-eventchange"),
- LSM_SCHEMA (sai_plat_owner_t, NULL, lsm_plat_list, "sai-builders"),
- LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
- /* shares struct */ "sai-overview"),
- LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo,
- /* shares struct */ "sai-tasklogs"),
-};
+extern const lws_struct_map_t lsm_schema_json_map[];
+extern size_t lsm_schema_json_map_array_size;
enum {
SAIS_WS_WEBSRV_RX_TASKCHANGE,
SAIS_WS_WEBSRV_RX_EVENTCHANGE,
SAIS_WS_WEBSRV_RX_SAI_BUILDERS,
SAIS_WS_WEBSRV_RX_OVERVIEW, /* deleted or added event */
- SAIS_WS_WEBSRV_RX_TASKLOGS /* new logs for task (ratelimited) */
+ SAIS_WS_WEBSRV_RX_TASKLOGS, /* new logs for task (ratelimited) */
+ SAIS_WS_WEBSRV_RX_LOADREPORT
};
+/*
+ * sai-web is receiving from sai-server
+ *
+ * This may come in chunks and is statefully parsed
+ * so it's not directly sensitive to size or fragmentation
+ */
static int
saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
@@ -72,15 +66,15 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
sai_browse_rx_evinfo_t *ei;
int n;
-// lwsl_user("%s: len %d, flags: %d\n", __func__, (int)len, flags);
+// lwsl_user("%s: RX from server -> sai-web: len %d, flags: %d\n", __func__, (int)len, flags);
// lwsl_hexdump_notice(buf, len);
if (flags & LWSSS_FLAG_SOM) {
memset(&m->a, 0, sizeof(m->a));
m->a.map_st[0] = lsm_schema_json_map;
- m->a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_json_map);
- m->a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_json_map);
- m->a.ac_block_size = 128;
+ m->a.map_entries_st[0] = lsm_schema_json_map_array_size;
+ m->a.map_entries_st[1] = lsm_schema_json_map_array_size;
+ m->a.ac_block_size = 4096;
lws_struct_json_init_parse(&m->ctx, NULL, &m->a);
}
@@ -88,16 +82,17 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
n = lejp_parse(&m->ctx, (uint8_t *)buf, (int)len);
if (n < LEJP_CONTINUE || (n >= 0 && !m->a.dest)) {
lwsac_free(&m->a.ac);
- lwsl_hexdump_notice(buf, len);
lwsl_notice("%s: srv->web JSON decode failed '%s'\n",
__func__, lejp_error_to_string(n));
+ lwsl_hexdump_notice(buf, len);
+
return LWSSSSRET_DISCONNECT_ME;
}
if (!(flags & LWSSS_FLAG_EOM))
return 0;
-// lwsl_notice("%s: schema idx %d\n", __func__, m->a.top_schema_index);
+ lwsl_notice("%s: schema idx %d parsed correctly from sai-server\n", __func__, m->a.top_schema_index);
switch (m->a.top_schema_index) {
@@ -115,19 +110,14 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
break;
case SAIS_WS_WEBSRV_RX_SAI_BUILDERS:
- lwsl_notice("%s: updated sai builder list\n", __func__);
+ lwsl_notice("%s: updated sai builder list (%d browsers)\n", __func__, vhd->browsers.count);
if (vhd->builders)
lwsac_detach(&vhd->builders);
vhd->builders = m->a.ac;
m->a.ac = NULL;
vhd->builders_owner =
&((sai_plat_owner_t *)m->a.dest)->plat_owner;
- lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) {
- struct pss *pss = lws_container_of(p, struct pss, same);
-
- saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY);
-
- } lws_end_foreach_dll(p);
+ saiw_ws_broadcast_raw(vhd, buf, len);
break;
case SAIS_WS_WEBSRV_RX_OVERVIEW:
@@ -155,6 +145,13 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
lws_callback_on_writable(pss->wsi);
} lws_end_foreach_dll(p);
break;
+
+ case SAIS_WS_WEBSRV_RX_LOADREPORT:
+ /* A builder sent a load report, forward to all browsers */
+ lwsl_notice("%s: ===== Received load report, broadcasting to %d browsers\n",
+ __func__, (int)vhd->browsers.count);
+ saiw_ws_broadcast_raw(vhd, buf, len);
+ break;
}
if (flags & LWSSS_FLAG_EOM)
diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c
index 1f1f42d..18995f9 100644
--- a/src/web/w-ws-browser.c
+++ b/src/web/w-ws-browser.c
@@ -31,6 +31,26 @@
#include "w-private.h"
/*
+ * This allows other parts of sai-web to queue a raw buffer to be sent to
+ * all connected browsers, eg, for load reports.
+ */
+void
+saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len)
+{
+ // 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);
+
+ if (lws_buflist_append_segment(&pss->raw_tx, buf, len) >= 0)
+ lws_callback_on_writable(pss->wsi);
+ } lws_end_foreach_dll(p);
+}
+
+extern const lws_struct_map_t lsm_load_report_members[2];
+
+/*
* For decoding specific event data request from browser
*/
@@ -66,6 +86,8 @@ static const lws_struct_map_t lsm_schema_json_map_bwsrx[] = {
/* shares struct */ "com.warmcat.sai.eventdelete"),
LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel,
"com.warmcat.sai.taskcan"),
+ LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members,
+ "com.warmcat.sai.loadreport"),
};
enum {
@@ -74,7 +96,7 @@ enum {
SAIM_WS_BROWSER_RX_TASKRESET,
SAIM_WS_BROWSER_RX_EVENTRESET,
SAIM_WS_BROWSER_RX_EVENTDELETE,
- SAIM_WS_BROWSER_RX_TASKCANCEL
+ SAIM_WS_BROWSER_RX_TASKCANCEL,
};
@@ -220,6 +242,7 @@ saiw_pss_schedule_eventinfo(struct pss *pss, const char *event_uuid)
goto bail;
sch->ov_db_done = 1;
+ // lwsl_warn("%s: doing WSS_PREPARE_BUILDER_SUMMARY\n", __func__);
saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY);
return 0;
@@ -329,6 +352,8 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub)
sch->logsub = !!logsub;
sch->one_event = lws_container_of(o.head, sai_event_t, list);
+ // lwsl_warn("%s: doing WSS_PREPARE_BUILDER_SUMMARY\n", __func__);
+
saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY);
return 0;
@@ -433,6 +458,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf,
/*
* he's asking for the overview schema
*/
+ // lwsl_warn("%s: SAIM_WS_BROWSER_RX_TASKINFO: doing WSS_PREPARE_BUILDER_SUMMARY\n", __func__);
saiw_alloc_sched(pss, WSS_PREPARE_OVERVIEW);
saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY);
@@ -583,6 +609,31 @@ again:
// lwsl_notice("%s: send_state %d, pss %p, wsi %p\n", __func__,
// pss->send_state, pss, pss->wsi);
+ if (pss->raw_tx) {
+ char som, eom;
+ int used;
+
+ p = start; /* buf + LWS_PRE */
+ used = lws_buflist_fragment_use(&pss->raw_tx, p,
+ lws_ptr_diff_size_t(end, p), &som, &eom);
+ if (!used)
+ return 0;
+
+ flags = lws_write_ws_flags(LWS_WRITE_TEXT, som, eom);
+ if (lws_write(pss->wsi, p, (size_t)used, (enum lws_write_protocol)flags) < 0)
+ return -1;
+
+ /*
+ * if there are more fragments, we must exit now and wait for
+ * the next writable callback to send the rest. Otherwise, we
+ * can fall through and check for other work to do.
+ */
+ if (pss->raw_tx) {
+ lws_callback_on_writable(pss->wsi);
+ return 0;
+ }
+ }
+
if (pss->sched.count)
sch = lws_container_of(pss->sched.head, saiw_scheduled_t, list);
else
@@ -1200,3 +1251,63 @@ no_sch:
return 0;
}
+
+/*
+ * This should be called from the browser-facing websocket protocol handler
+ * on LWS_CALLBACK_ESTABLISHED and LWS_CALLBACK_CLOSED events to keep an
+ * accurate real-time list of connected browsers.
+ */
+void
+saiw_browser_state_changed(struct pss *pss, int established)
+{
+ if (established)
+ lws_dll2_add_tail(&pss->same, &pss->vhd->browsers);
+ else
+ lws_dll2_remove(&pss->same);
+
+ /*
+ * After any change, recalculate the total and inform the server
+ */
+ saiw_update_viewer_count(pss->vhd);
+}
+
+/*
+ * This function calculates the current number of connected browsers and
+ * sends an update to the sai-server.
+ */
+void
+saiw_update_viewer_count(struct vhd *vhd)
+{
+ sai_viewer_state_t vs;
+ char buf[LWS_PRE + 256];
+ size_t len;
+
+ if (!vhd || !vhd->h_ss_websrv)
+ return;
+
+ /* The count is simply the number of items in the browsers list */
+ vs.viewers = (unsigned int)vhd->browsers.count;
+
+ const lws_struct_map_t lsm_viewercount_members[] = {
+ LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"),
+ };
+
+ const lws_struct_map_t lsm_schema_json_map[] = {
+ LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewercount_members,
+ "com.warmcat.sai.viewercount"),
+ };
+
+ lws_struct_serialize_t *js = lws_struct_json_serialize_create(
+ lsm_schema_json_map, LWS_ARRAY_SIZE(lsm_schema_json_map),
+ 0, &vs);
+ if (!js)
+ return;
+
+ len = 0;
+ lws_struct_json_serialize(js, (unsigned char *)buf + LWS_PRE,
+ sizeof(buf) - LWS_PRE, &len);
+ lws_struct_json_serialize_destroy(&js);
+
+ if (len > 0)
+ saiw_websrv_queue_tx(vhd->h_ss_websrv, buf + LWS_PRE, len);
+}