Author: Andy Green Date: Fri Aug 01 14:46:42 2025 +0100 load reporting lr1 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 = `
`; + innerHTML += ``; + innerHTML += ``; + innerHTML += ``; + innerHTML += `
${plat.peer_ip}`; + innerHTML += `
`; + + for (let i = 0; i < plat.instances; i++) { + innerHTML += `
`; + } + + innerHTML += `
`; + 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 += "
" + + (e.peer_ip ? " / " + san(e.peer_ip) : "") + + "\" id=\"binfo-" + san(e.name) + "\">" + + "
" + sai_plat_icon(e.platform, 1) + (e.peer_ip ? "
" + san(e.peer_ip) : ""); + /* Add a container for the instance load boxes */ + s += "
"; + for (var i = 0; i < e.instances; i++) { + s += "
"; + } + s += "
"; + 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 += "
 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 #if !defined(WIN32) #include +#else +#include #endif #include #include @@ -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 #include +#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 #include +#if !defined(WIN32) #include #include +#endif #if defined(__linux__) #include 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); +}