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;