Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
Author[]Andy Green <andy@warmcat.com> 2025-08-08 16:45 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-08-10 10:47 UTC
Treeca466c0930b60a7a9d94de750a3218705b327598   Raw Patch
 
debug-delta
debug-delta
diff --git a/assets/sai.css b/assets/sai.css index 05fd1d7..5ea141b 100644 --- a/assets/sai.css +++ b/assets/sai.css @@ -689,8 +689,12 @@ div.awaiting { div.bdr { display: inline-block; opacity: 1; + transition: opacity 2s ease; } +.bdr.offline { + opacity: 0.35; +} img.branch { diff --git a/assets/sai.js b/assets/sai.js index 49d5411..4882e2b 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -845,11 +845,33 @@ function getBuilderGroupKey(platName) { return hostname; } +function refresh_state(task_uuid, task_state) +{ + var tsi = document.getElementById("taskstate_" + task_uuid); + + if (tsi) { + tsi.classList.remove("taskstate1"); + tsi.classList.remove("taskstate2"); + tsi.classList.remove("taskstate3"); + tsi.classList.remove("taskstate4"); + tsi.classList.remove("taskstate5"); + tsi.classList.remove("taskstate6"); + tsi.classList.remove("taskstate7"); + tsi.classList.add("taskstate" + task_state); + // console.log("refresh_state taskstate" + task_state); + } +} +function after_delete() { + location.reload(); +} function createBuilderDiv(plat) { const platDiv = document.createElement("div"); platDiv.className = "ibuil bdr"; + if (!plat.online) + platDiv.className += " offline"; + platDiv.id = "binfo-" + plat.name; platDiv.title = plat.platform + "@" + plat.name.split('.')[0] + " / " + plat.peer_ip; @@ -863,7 +885,7 @@ function createBuilderDiv(plat) { 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 class="instload" id="instload-${plat.name}">`; // Changed class name for clarity + innerHTML += `<div class="instload" id="instload-${plat.name}">`; // Create initial idle squares for (let i = 0; i < plat.instances; i++) { @@ -878,112 +900,6 @@ function createBuilderDiv(plat) { return platDiv; } -function render_builders(jso) -{ - var s, n, conts = [], nc = 0; - - s = "<table class=\"builders\"><tr><td>"; - -// s = "<table class=\"builders\"><tr><td><img class=\"ip bsvg\" " + -// "src=\"/sai/builder.png\"></td><td>"; - - for (n = 0; n < jso.builders.length; n++) { - var e = jso.builders[n], host; - - host = e.name.split('.')[0]; - if (host.split('-')[1]) { - var m; - - for (m = 0; m < conts.length; m++) - if (host.split('-')[1] === conts[m]) - m = 999; - if (m < 999 || !conts.length) - conts.push(host.split('-')[1]); - } - } - - for (nc = 0; nc <= conts.length; nc++) { - var samplat = ""; - - /* - * nc == conts.length means those not - * in a container - */ - - if (nc < conts.length) - s += "<div class=\"ibuil ibuilctr bdr\"><div class=\"ibuilctrname bdr\">" + conts[nc] + "</div>"; - - for (n = 0; n < jso.builders.length; n++) { - var e = jso.builders[n], nn, host, plat, cia, sen; - - sen = e.name; - host = sen.split('.')[0]; - cia = host.split('-')[1]; - - // console.log("host " + host + ", cia " + cia); - - if ((nc == conts.length && !cia) || - (cia && cia === conts[nc])) { - var did = 0, arc; - - if (!n || - host !== jso.builders[n - 1].name.split('.')[0] || - e.platform != samplat) { - - s += "<div class=\"ibuil bdr\" title=\"" + - san(e.platform) + "@" + san(host) + - (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)/*.split('.')[0]*/ + "\">"; - for (var i = 0; i < e.instances; i++) { - s += "<div class=\"inst_box inst_idle\" title=\"instance " + i + ": idle\"><div class=\"inst_bar\"></div></div>"; - } - s += "</div>"; - - samplat = e.platform; - did = 1; - } - - if (n + 1 == jso.builders.length || did) - s += "</td></tr></table></div>"; - } - } - if (nc < conts.length) - s += "</div>"; - } - - s += "</td></tr></table></td></tr>"; - s = s + "</table>"; - - return s; -} - -function refresh_state(task_uuid, task_state) -{ - var tsi = document.getElementById("taskstate_" + task_uuid); - - if (tsi) { - tsi.classList.remove("taskstate1"); - tsi.classList.remove("taskstate2"); - tsi.classList.remove("taskstate3"); - tsi.classList.remove("taskstate4"); - tsi.classList.remove("taskstate5"); - tsi.classList.remove("taskstate6"); - tsi.classList.remove("taskstate7"); - tsi.classList.add("taskstate" + task_state); - // console.log("refresh_state taskstate" + task_state); - } -} - -function after_delete() { - location.reload(); -} - function ws_open_sai() { var s = "", q, qa, qi, q5, q5s; @@ -1136,16 +1052,97 @@ function ws_open_sai() return; // Stop processing this old message } - switch (jso.schema) { + switch (jso.schema) { - case "com.warmcat.sai.builders": + case "com.warmcat.sai.builders": + const buildersContainer = document.getElementById("sai_builders"); + if (!buildersContainer) { break; } + + let platformsArray = null; + if (jso.platforms && Array.isArray(jso.platforms)) { + platformsArray = jso.platforms; + } else if (jso.builders && Array.isArray(jso.builders)) { + platformsArray = jso.builders; + } + if (!platformsArray) { + buildersContainer.innerHTML = ""; // Clear display if data is invalid + break; + } + + // --- Reconciliation Logic --- + + // Step 1: Find the main content area (the TD). Create it if this is the first run. + let tdContainer = buildersContainer.querySelector("table td"); + if (!tdContainer) { + buildersContainer.innerHTML = ""; // Clear for safety + const table = document.createElement("table"); + table.className = "builders"; + const tbody = document.createElement("tbody"); + const tr = document.createElement("tr"); + tdContainer = document.createElement("td"); + tr.appendChild(tdContainer); + tbody.appendChild(tr); + table.appendChild(tbody); + buildersContainer.appendChild(table); + } + + // Step 2: Detach all existing builder divs and store them in a map for reuse. + // This preserves them so their CSS transitions will work. + const existingDivs = new Map(); + tdContainer.querySelectorAll(".bdr").forEach(div => { + const name = div.id.substring(6); // "binfo-" is 6 chars + if (name) { + existingDivs.set(name, div); + } + }); + tdContainer.innerHTML = ""; // Clear the container, but the divs are still in memory. + + // Step 3: Rebuild the group structure from scratch, reusing the old divs. + platformsArray.sort((a, b) => { + if (a.name && b.name) { + return a.name.localeCompare(b.name); + } + return 0; // Don't sort if names are missing + }); + + const platformsByGroup = {}; + for (const plat of platformsArray) { + const groupKey = getBuilderGroupKey(plat.name); + if (!platformsByGroup[groupKey]) platformsByGroup[groupKey] = []; + platformsByGroup[groupKey].push(plat); + } + + const groupKeys = Object.keys(platformsByGroup).sort(); + for (const key of groupKeys) { + const groupPlatforms = platformsByGroup[key]; + const nestedPlatforms = groupPlatforms.filter(p => getBuilderHostname(p.name) !== key); + const mainPlatforms = groupPlatforms.filter(p => getBuilderHostname(p.name) === key); + + 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); + + nestedPlatforms.forEach(plat => { + const div = existingDivs.get(plat.name) || createBuilderDiv(plat); + div.classList.toggle('offline', !plat.online); + groupDiv.appendChild(div); + }); + tdContainer.appendChild(groupDiv); + } + + mainPlatforms.forEach(plat => { + const div = existingDivs.get(plat.name) || createBuilderDiv(plat); + div.classList.toggle('offline', !plat.online); + tdContainer.appendChild(div); + }); + } + break; - s = render_builders(jso); - - if (document.getElementById("sai_builders")) - document.getElementById("sai_builders").innerHTML = s; - break; - case "sai.warmcat.com.overview": /* * Sent with an array of e[] to start, but also @@ -1424,82 +1421,6 @@ function ws_open_sai() } 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; - } - - jso.platforms.sort((a, b) => { - if (a.name && b.name) { - return a.name.localeCompare(b.name); - } - return 0; // Don't sort if names are missing - }); - - // 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]; - - // 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.platforms || !Array.isArray(jso.platforms)) { break; diff --git a/src/common/include/private.h b/src/common/include/private.h index ae50033..0720c6b 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -460,7 +460,8 @@ extern const lws_struct_map_t lsm_schema_json_map_event[1], lsm_resource[4] ; -extern const lws_struct_map_t lsm_plat[8]; +extern const lws_struct_map_t lsm_plat[6]; +extern const lws_struct_map_t lsm_plat_for_json[8]; extern const lws_ss_info_t ssi_said_logproxy; extern struct lws_ss_handle *ssh[3]; @@ -478,3 +479,5 @@ saicom_lp_callback_on_drain(saicom_drain_cb cb, void *opaque); void sul_idle_cb(lws_sorted_usec_list_t *sul); + + diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 3fb39fc..e612330 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -41,23 +41,33 @@ const lws_struct_map_t lsm_load_report_members[] = { }; const lws_struct_map_t lsm_plat[] = { /* !!! keep extern length in common/include/private.h in sync */ - LSM_UNSIGNED (sai_event_t, uid, "uid"), + LSM_UNSIGNED (sai_plat_t, uid, "uid"), LSM_STRING_PTR (sai_plat_t, name, "name"), - LSM_UNSIGNED (sai_plat_t, ongoing, "ongoing"), LSM_UNSIGNED (sai_plat_t, instances, "instances"), LSM_STRING_PTR (sai_plat_t, platform, "platform"), - LSM_SIGNED (sai_plat_t, online, "online"), LSM_UNSIGNED (sai_plat_t, last_seen, "last_seen"), LSM_CARRAY (sai_plat_t, peer_ip, "peer_ip"), }; +// This is the map for serializing to JSON +const lws_struct_map_t lsm_plat_for_json[] = { + LSM_UNSIGNED(sai_plat_t, uid, "uid"), + LSM_STRING_PTR(sai_plat_t, name, "name"), + LSM_SIGNED(sai_plat_t, ongoing, "ongoing"), // MUST be present + LSM_SIGNED(sai_plat_t, instances, "instances"), + LSM_STRING_PTR(sai_plat_t, platform,"platform"), + LSM_SIGNED(sai_plat_t, online, "online"), // MUST be present + LSM_UNSIGNED(sai_plat_t, last_seen, "last_seen"), + LSM_CARRAY(sai_plat_t, peer_ip, "peer_ip"), +}; + const lws_struct_map_t lsm_schema_map_plat_simple[] = { - LSM_SCHEMA (sai_plat_t, NULL, lsm_plat, "com-warmcat-sai-ba"), + LSM_SCHEMA (sai_plat_t, NULL, lsm_plat_for_json, "com-warmcat-sai-ba"), }; const lws_struct_map_t lsm_plat_list[] = { LSM_LIST (sai_plat_owner_t, plat_owner, sai_plat_t, - sai_plat_list, NULL, lsm_plat, "platforms"), + sai_plat_list, NULL, lsm_plat_for_json, "builders"), }; const lws_struct_map_t lsm_schema_map_plat[] = { diff --git a/src/server/s-comms.c b/src/server/s-comms.c index e069cad..85aed2e 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -508,7 +508,9 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, sai_sqlite3_statement(vhd->server.pdb, "CREATE UNIQUE INDEX IF NOT EXISTS name_idx ON builders (name)", "create builder name index"); - + +// sais_mark_all_builders_offline(vhd); + lwsl_notice("%s: creating server stream\n", __func__); if (lws_ss_create(vhd->context, 0, &ssi_server, vhd, @@ -759,6 +761,8 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, /* remove pss from vhd->builders (active connection list) */ lws_dll2_remove(&pss->same); + sais_builder_disconnected(vhd, wsi); + /* * Find any builder-tracking objects that were using this departing * connection. Mark them as offline in the database. diff --git a/src/server/s-private.h b/src/server/s-private.h index 1098f26..2c72ba7 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -315,3 +315,13 @@ sais_resource_rr_destroy(sai_resource_requisition_t *rr); int sais_platforms_with_tasks_pending(struct vhd *vhd); + +sai_plat_t * +sais_builder_from_uuid(struct vhd *vhd, const char *hostname, const char *_file, int _line); + +void +sais_builder_disconnected(struct vhd *vhd, struct lws *wsi); + + +void +sais_mark_all_builders_offline(struct vhd *vhd); diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index 2d351bb..b51f541 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -79,6 +79,21 @@ enum { SAIS_WS_WEBSRV_RX_VIEWERCOUNT, }; +void +sais_mark_all_builders_offline(struct vhd *vhd) +{ + char *err = NULL; + + lwsl_notice("%s: marking all builders offline initially\n", __func__); + + sqlite3_exec(vhd->server.pdb, "UPDATE builders SET online = 0;", + NULL, NULL, &err); + if (err) { + lwsl_err("%s: sqlite error: %s\n", __func__, err); + sqlite3_free(err); + } +} + int sais_validate_id(const char *id, int reqlen) { @@ -149,68 +164,74 @@ 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) { - lws_dll2_owner_t dbo; + lws_dll2_owner_t db_builders_owner; struct lwsac *ac = NULL; char *p = vhd->json_builders, *end = p + sizeof(vhd->json_builders), subsequent = 0; lws_struct_serialize_t *js; - sai_plat_t *b; + sai_plat_t *builder_from_db; size_t w; - int n; - /* - * Query the database for ALL builders, online and offline, - * sorted by name. - */ + memset(&db_builders_owner, 0, sizeof(db_builders_owner)); + if (lws_struct_sq3_deserialize(vhd->server.pdb, NULL, "name ", - lsm_schema_sq3_map_plat, &dbo, &ac, 0, 100)) { + lsm_schema_sq3_map_plat, + &db_builders_owner, &ac, 0, 100)) { lwsl_err("%s: Failed to query builders from DB\n", __func__); return 1; } + lwsl_warn("%s: count deserialized %d\n", __func__, (int)db_builders_owner.count); + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), - "{\"schema\":\"sai-builders\"," - "\"platforms\":["); + "{\"schema\":\"com.warmcat.sai.builders\",\"builders\":["); + + lws_start_foreach_dll(struct lws_dll2 *, walk, db_builders_owner.head) { + sai_plat_t *live_builder; - lws_start_foreach_dll(struct lws_dll2 *, walk, dbo.head) { + builder_from_db = lws_container_of(walk, sai_plat_t, sai_plat_list); - b = lws_container_of(walk, sai_plat_t, sai_plat_list); + /* + * Find this builder in the live list by name. This is safe because + * builder_from_db->name is a valid string within the scope of this function. + */ + live_builder = sais_builder_from_uuid(vhd, builder_from_db->name, __FILE__, __LINE__); + + if (live_builder) { + builder_from_db->online = 1; + builder_from_db->ongoing = live_builder->ongoing; + lws_strncpy(builder_from_db->peer_ip, live_builder->peer_ip, + sizeof(builder_from_db->peer_ip)); + } else { + builder_from_db->online = 0; + builder_from_db->ongoing = 0; + } js = lws_struct_json_serialize_create( lsm_schema_map_plat_simple, LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), - 0, b); + 0, builder_from_db); if (!js) { - lwsl_err("%s: json serialize create failed\n", __func__); goto bail; } if (subsequent) *p++ = ','; subsequent = 1; - n = (int)lws_struct_json_serialize(js, (unsigned char *)p, - lws_ptr_diff_size_t(end, p), &w); - p += w; - lws_struct_json_serialize_destroy(&js); - - if (n == LSJS_RESULT_ERROR) { - lwsl_err("%s: json serialize failed\n", __func__); + if (lws_struct_json_serialize(js, (uint8_t *)p, + lws_ptr_diff_size_t(end, p), &w) != LSJS_RESULT_FINISH) { + lws_struct_json_serialize_destroy(&js); goto bail; } + p += w; + lws_struct_json_serialize_destroy(&js); } lws_end_foreach_dll(walk); - /* end of the list of builders */ - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - /* - * This is the SERVER's WEB daemon server, broadcasting to all connected - * clients (the WEB daemons)... the list of BUILDERS - */ sais_websrv_broadcast(vhd->h_ss_websrv, vhd->json_builders, lws_ptr_diff_size_t(p, vhd->json_builders)); @@ -550,6 +571,8 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, *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); + if (m->bltx) return lws_ss_request_tx(m->ss); @@ -589,7 +612,7 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, struct pss *pss_builder = lws_container_of(p, struct pss, same); sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); if (vsend) { - vsend->viewers = new_viewers_present; + vsend->viewers = (unsigned int)new_viewers_present; lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); lws_callback_on_writable(pss_builder->wsi); } diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 9de422e..69b29df 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -239,22 +239,52 @@ sais_log_to_db(struct vhd *vhd, sai_log_t *log) sais_dump_logs_to_db, 250 * LWS_US_PER_MS); } -static sai_plat_t * -sais_builder_from_uuid(struct vhd *vhd, const char *hostname) +sai_plat_t * +sais_builder_from_uuid(struct vhd *vhd, const char *hostname, const char *_file, int _line) { lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { sai_plat_t *cb = lws_container_of(p, sai_plat_t, sai_plat_list); - if (!strcmp(hostname, cb->name)) + if (!strcmp(hostname, cb->name)) { + lwsl_err("%s: %s:%d: found live builder %s\n", __func__, _file, _line, hostname); + cb->online = 1; return cb; + } } lws_end_foreach_dll(p); return NULL; } +/* + * Called from the builder protocol LWS_CALLBACK_CLOSED handler + */ +void +sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) +{ + sai_plat_t *cb; + + /* + * A builder's websocket has closed. Find all platforms associated + * with it, mark them as offline in the database, and remove them + * from the live in-memory list. + */ + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + vhd->server.builder_owner.head) { + cb = lws_container_of(p, sai_plat_t, sai_plat_list); + + if (cb->wsi == wsi) { + lwsl_notice("%s: Builder '%s' disconnected\n", __func__, + cb->name); + + lws_dll2_remove(&cb->sai_plat_list); + free(cb); + } + } lws_end_foreach_dll_safe(p, p1); +} + int sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name) { @@ -355,8 +385,6 @@ handle: build = lws_container_of(pb, sai_plat_t, sai_plat_list); sai_plat_t *live_cb; - lwsl_notice("%s: seeing plat %s\n", __func__, build->name); - /* * Step 1: Upsert this platform into the persistent database. */ @@ -372,16 +400,24 @@ handle: /* * Step 2: Update the long-lived, malloc'd in-memory list. */ - live_cb = sais_builder_from_uuid(vhd, build->name); + //cb = sais_builder_from_uuid(vhd, build->name); + //if (cb) + // sais_builder_disconnected(vhd, cb->wsi); + live_cb = sais_builder_from_uuid(vhd, build->name, __FILE__, __LINE__); if (live_cb) { /* Already exists (reconnect), just update dynamic info */ + lwsl_err("%s: found live builder for %s\n", __func__, build->name); live_cb->wsi = pss->wsi; live_cb->ongoing = 0; /* Reset ongoing task count on connect */ lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip)); + live_cb->online = 1; } else { /* New builder, create a deep-copied, malloc'd object */ size_t nlen = strlen(build->name) + 1; size_t plen = strlen(build->platform) + 1; + + lwsl_err("%s: no live for %s\n", __func__, build->name); + live_cb = malloc(sizeof(*live_cb) + nlen + plen); if (live_cb) { char *p_str = (char *)(live_cb + 1); @@ -392,6 +428,7 @@ handle: memcpy(p_str + nlen, build->platform, plen); live_cb->instances = build->instances; live_cb->wsi = pss->wsi; + live_cb->online = 1; lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip)); lws_dll2_add_tail(&live_cb->sai_plat_list, &vhd->server.builder_owner); } @@ -489,7 +526,7 @@ bail: rej = (sai_rejection_t *)pss->a.dest; rej->host_platform[sizeof(rej->host_platform) - 1] = '\0'; - cb = sais_builder_from_uuid(vhd, rej->host_platform); + cb = sais_builder_from_uuid(vhd, rej->host_platform, __FILE__, __LINE__); if (!cb) { lwsl_info("%s: unknown builder %s rejecting\n", __func__, rej->host_platform); diff --git a/src/web/w-comms.c b/src/web/w-comms.c index 87cb552..1997a49 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -92,7 +92,7 @@ const lws_struct_map_t lsm_schema_json_map[] = { /* 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_plat_owner_t, NULL, lsm_plat_list, "com.warmcat.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, diff --git a/src/web/w-private.h b/src/web/w-private.h index 97281bb..38d5a10 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -205,13 +205,13 @@ typedef struct sais_sqlite_cache { } sais_sqlite_cache_t; struct vhd { - struct lws_context *context; - struct lws_vhost *vhost; + struct lws_context *context; + struct lws_vhost *vhost; /* pss lists */ - struct lws_dll2_owner browsers; + struct lws_dll2_owner browsers; - struct lws_dll2_owner *builders_owner; + struct lws_dll2_owner builders_owner; struct lwsac *builders; /* our keys */ diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c index 2de152c..ccc96c1 100644 --- a/src/web/w-websrv.c +++ b/src/web/w-websrv.c @@ -34,10 +34,7 @@ typedef struct saiw_websrv { lws_struct_args_t a; struct lejp_ctx ctx; - //lws_dll2_t struct lws_buflist *bltx; - struct lwsac *deprecated; - } saiw_websrv_t; extern const lws_struct_map_t lsm_schema_json_map[]; @@ -47,9 +44,9 @@ 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_LOADREPORT + SAIS_WS_WEBSRV_RX_OVERVIEW, /* deleted or added event */ + SAIS_WS_WEBSRV_RX_TASKLOGS, /* new logs for task (ratelimited) */ + SAIS_WS_WEBSRV_RX_LOADREPORT, /* builder's cpu load report */ }; /* @@ -58,7 +55,6 @@ enum { * 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) { @@ -67,41 +63,64 @@ 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: RX from server -> sai-web: len %d, flags: %d\n", __func__, (int)len, flags); -// lwsl_hexdump_notice(buf, len); + // lwsl_warn("%s: len %d, flags %d\n", __func__, (int)len, flags); if (flags & LWSSS_FLAG_SOM) { - m->deprecated = vhd->builders; + /* First fragment of a new message. Clear old parse results and init. */ + lwsac_free(&m->a.ac); memset(&m->a, 0, sizeof(m->a)); m->a.map_st[0] = lsm_schema_json_map; 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); } + // fprintf(stderr, "%s: rx: %.*s\n", __func__, (int)len, buf); + n = lejp_parse(&m->ctx, (uint8_t *)buf, (int)len); - if (n < LEJP_CONTINUE || (n >= 0 && !m->a.dest)) { - vhd->builders_owner = NULL; - lwsac_free(&m->a.ac); + + /* Check for fatal error OR completion without an object */ + if (n < 0 && n != LEJP_CONTINUE) { 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; + goto cleanup_and_disconnect; } -// if (!(flags & LWSSS_FLAG_EOM)) -// return 0; + /* + * This is the key: if the message is not yet complete, just return + * and wait for the next fragment. Don't process anything yet. + */ + if (n == LEJP_CONTINUE) { + /* + * Also forward this fragment to browsers if the message is for them. + * We can check the schema index which is available after the + * "schema" member is parsed, even on the first fragment. + */ + switch (m->a.top_schema_index) { + case SAIS_WS_WEBSRV_RX_LOADREPORT: + saiw_ws_broadcast_raw(vhd, buf, len, 0, + lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); + break; + } - // lwsl_notice("%s: schema idx %d parsed correctly from sai-server\n", __func__, m->a.top_schema_index); + return 0; + } + + /* + * If we get here, the message is fully parsed (n >= 0). + * Now we can safely process m->a.dest. + */ + if (!m->a.dest) { + lwsl_warn("%s: JSON parsed but produced no object\n", __func__); + goto cleanup_parse_allocs; + } switch (m->a.top_schema_index) { case SAIS_WS_WEBSRV_RX_TASKCHANGE: ei = (sai_browse_rx_evinfo_t *)m->a.dest; - /* server has told us of a task change */ lwsl_notice("%s: TASKCHANGE %s\n", __func__, ei->event_hash); saiw_browsers_task_state_change(vhd, ei->event_hash); break; @@ -113,65 +132,67 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) break; case SAIS_WS_WEBSRV_RX_SAI_BUILDERS: + lwsac_free(&vhd->builders); + lws_dll2_owner_clear(&vhd->builders_owner); + vhd->builders = m->a.ac; + m->a.ac = NULL; /* The vhd now owns this memory */ - /* vhd holds a pointer to the active ac and a pointer to the owner (also lives in the ac) */ + /* Move the parsed objects to the vhd's list */ + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + ((sai_plat_owner_t *)m->a.dest)->plat_owner.head) { + sai_plat_t *cb = lws_container_of(p, sai_plat_t, sai_plat_list); - // lwsl_notice("%s: updated sai builder list\n", __func__); - if (vhd->builders) - lwsac_detach(&vhd->builders); + lws_dll2_remove(&cb->sai_plat_list); + lws_dll2_add_tail(&cb->sai_plat_list, &vhd->builders_owner); + } lws_end_foreach_dll_safe(p, p1); - /* we take over ownership of the ac */ + /* schedule emitting the builder summary to each browser */ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { + struct pss *pss = lws_container_of(p, struct pss, same); - vhd->builders = m->a.ac; - m->a.ac = NULL; - vhd->builders_owner = &((sai_plat_owner_t *)m->a.dest)->plat_owner; - if (lwsac_assert_valid(vhd->builders, vhd->builders_owner, sizeof(lws_dll2_owner_t))) - break; - saiw_ws_broadcast_raw(vhd, buf, len, 0, - lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); + saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); + } lws_end_foreach_dll(p); break; case SAIS_WS_WEBSRV_RX_OVERVIEW: lwsl_notice("%s: force overview\n", __func__); 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_OVERVIEW); } lws_end_foreach_dll(p); break; case SAIS_WS_WEBSRV_RX_TASKLOGS: ei = (sai_browse_rx_evinfo_t *)m->a.dest; - /* - * ratelimited indication that logs for a particular task - * changed... for each connected browser subscribed to logs for - * that task, let them know - */ - lws_start_foreach_dll(struct lws_dll2 *, p, - vhd->subs_owner.head) { - struct pss *pss = lws_container_of(p, struct pss, - subs_list); - + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->subs_owner.head) { + struct pss *pss = lws_container_of(p, struct pss, subs_list); if (!strcmp(pss->sub_task_uuid, ei->event_hash)) 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, 2, - lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); + /* Forward the final fragment of the load report */ + saiw_ws_broadcast_raw(vhd, buf, len, 2, + lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); break; } -// if (flags & LWSSS_FLAG_EOM && m->deprecated) -// lwsac_free(&m->deprecated); - +cleanup_parse_allocs: + /* + * Free the memory used for THIS parse. + * In the BUILDERS case, m->a.ac was transferred to vhd->builders, + * so it will be NULL here and lwsac_free is a no-op. + */ + lwsac_free(&m->a.ac); return 0; + +cleanup_and_disconnect: + lwsac_free(&m->a.ac); + return LWSSSSRET_DISCONNECT_ME; } + static int saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, int *flags) diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 63fc4b2..96bb119 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -1030,15 +1030,12 @@ so_finish: lws_json_purify(esc1, pss->auth_user, sizeof(esc1) - 1, &iu)); if (vhd && vhd->builders) { - lwsac_reference(vhd->builders); - sch->walk = lws_dll2_get_head(vhd->builders_owner); + // lwsac_reference(vhd->builders); + sch->walk = lws_dll2_get_head(&vhd->builders_owner); - /* builders_owner must be inside vhd->builders ac */ - if (lwsac_assert_valid(vhd->builders, vhd->builders_owner, sizeof(lws_dll2_owner_t))) - break; - /* HEAD of the owner list must be also inside the vhd->builders ac */ - if (sch->walk && lwsac_assert_valid(vhd->builders, sch->walk, sizeof(sai_plat_t))) - break; + /* HEAD of the owner list must be inside the vhd->builders ac */ + // if (sch->walk && lwsac_assert_valid(vhd->builders, sch->walk, sizeof(sai_plat_t))) + // break; } else { lwsl_notice("%s: BUILDER_SUMMARY: can't start walk\n", __func__); sch->walk = 0; @@ -1065,20 +1062,16 @@ so_finish: * builders / platforms we feel are connected to us */ - lwsl_notice("%s: WSS_SEND_BUILDER_SUMMARY outside write loop, walk %p\n", __func__, sch->walk); - while (end - p > 512 && sch->walk && pss->send_state == WSS_SEND_BUILDER_SUMMARY) { - /* every builder must be also inside the vhd->builders ac */ - if (lwsac_assert_valid(vhd->builders, sch->walk, sizeof(sai_plat_t))) - break; + /* every builder must be inside the vhd->builders ac */ + //if (lwsac_assert_valid(vhd->builders, sch->walk, sizeof(sai_plat_t))) + // break; sai_plat_t *b = lws_container_of(sch->walk, sai_plat_t, sai_plat_list); - lwsl_notice("%s: serializing inside %s\n", __func__, b->name); - js = lws_struct_json_serialize_create( lsm_schema_map_plat_simple, LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), @@ -1118,7 +1111,7 @@ so_finish: break; b_finish: p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - lwsac_unreference(&vhd->builders); + // lwsac_unreference(&vhd->builders); endo = 1; break;
Page fetched 0s ago, creation time: 10ms (vhost etag hits: 0%, cache hits: 0%)