diff --git a/assets/sai.css b/assets/sai.css
index a3f9fd2..a8cf241 100644
--- a/assets/sai.css
+++ b/assets/sai.css
@@ -194,7 +194,7 @@ span.tty1 {
span.noscript {
- text-size:200%;
+ font-size: 200%;
}
span.nowrap {
@@ -371,6 +371,13 @@ img.bico {
transition: background-color 500ms linear;
}
+.taskstate10 {
+ background:#ffdc00;
+ color:#302000;
+ opacity:1;
+ transition: background-color 500ms linear;
+}
+
div.stats_hidden {
display: none;
}
@@ -472,7 +479,6 @@ div.sai_arts {
span.ti1 {
font-weight: normal;
font-size: 9pt;
- //background:#d0c0b0;
color:#000000;
padding:3px;
border-radius:3px;
@@ -656,6 +662,10 @@ div.ibuil {
width: 80px;
}
+div.ibuil table {
+ table-layout: fixed;
+}
+
div.ibuil.power-unmanaged {
border: 2px solid black;
}
@@ -841,6 +851,15 @@ img.branch {
background-color: #f2f2f2;
}
+.context-menu li.read-only {
+ cursor: default;
+ color: #888;
+}
+
+.context-menu li.read-only:hover {
+ background-color: inherit;
+}
+
.float-right {
float: right;
}
@@ -1088,8 +1107,8 @@ body.overlay-active {
.pcon-header {
font-weight: bold;
- font-size: 1.0em;
- padding: 4px;
+ font-size: 0.5em;
+ padding: 2px;
background-color: #e0e0e0;
border-radius: 3px;
display: flex;
@@ -1122,3 +1141,27 @@ body.overlay-active {
.pcon-orphans .pcon-header {
background-color: #ffe0e0;
}
+
+.builder-name-row {
+ display: flex;
+ justify-content: space-between;
+ align-items: center;
+ width: 100%;
+ margin-bottom: 2px;
+}
+
+.builder-short-name {
+ font-weight: bold;
+ overflow: hidden;
+ text-overflow: ellipsis;
+ white-space: nowrap;
+ flex-grow: 1;
+ text-align: left;
+ min-width: 0;
+}
+
+.builder-icons {
+ display: flex;
+ align-items: center;
+ flex-shrink: 0;
+}
diff --git a/assets/sai.js b/assets/sai.js
index 7889afc..58276b4 100644
--- a/assets/sai.js
+++ b/assets/sai.js
@@ -395,7 +395,7 @@ var lang_zhs = "{" +
"}}";
var logs = "", redpend = 0, gitohashi_integ = 0, authd = 0, exptimer, auth_user = "",
- logAnsiState = {},
+ logAnsiState = {}, logs_pending = "", lines_pending = "", times_pending = "",
ongoing_task_activities = {}, last_log_timestamp = 0, spreadsheet_data_cache = {}, loadreport_data_cache = {},
fadingTasks = new Map();
@@ -650,9 +650,13 @@ function updateTaskRow(tr, task, now_ut) {
s1 += "□";
qc++;
}
- tr.innerHTML = `<td>${s1}</td>` +
+ const newHTML = `<td>${s1}</td>` +
`<td>${agify(now_ut, task.started)} ago</td>` +
`<td><a href="index.html?task=${hsanitize(task.task_uuid)}">${hsanitize(task.task_name)}</a></td>`;
+
+ if (tr.innerHTML !== newHTML) {
+ tr.innerHTML = newHTML;
+ }
}
function updateSpreadsheetDOM(container, tasks) {
@@ -720,8 +724,11 @@ function updateSpreadsheetDOM(container, tasks) {
return (taskB.started - taskA.started) || taskA.task_name.localeCompare(taskB.task_name);
});
- for (const row of rows) {
- tbody.appendChild(row);
+ for (let i = 0; i < rows.length; i++) {
+ const expectedRow = rows[i];
+ if (tbody.children[i] !== expectedRow) {
+ tbody.insertBefore(expectedRow, tbody.children[i] || null);
+ }
}
}
@@ -774,6 +781,8 @@ function agify(now, secs)
i18n(age_names[n]) + "</span>";
}
+var aging_timer = null;
+
function aging()
{
var n, next = 24 * 3600,
@@ -812,7 +821,9 @@ function aging()
* Eg, if everything is counted in hours already, once per
* 5 minutes is accurate enough.
*/
- window.setTimeout(aging, next * 1000);
+ if (aging_timer)
+ clearTimeout(aging_timer);
+ aging_timer = window.setTimeout(aging, next * 1000);
}
var sai, jso, s, sai_arts = "";
@@ -1158,6 +1169,7 @@ function refresh_state(task_uuid, task_state)
tsi.classList.remove("taskstate5");
tsi.classList.remove("taskstate6");
tsi.classList.remove("taskstate7");
+ tsi.classList.remove("taskstate10");
tsi.classList.add("taskstate" + task_state);
// console.log("refresh_state taskstate" + task_state);
}
@@ -1173,7 +1185,8 @@ function createContextMenu(event, menuItems) {
// Remove any existing context menu
const existingMenus = document.querySelectorAll(".context-menu");
existingMenus.forEach(menu => {
- document.body.removeChild(menu);
+ if (document.body.contains(menu))
+ document.body.removeChild(menu);
});
const menu = document.createElement("div");
@@ -1184,26 +1197,51 @@ function createContextMenu(event, menuItems) {
const ul = document.createElement("ul");
menu.appendChild(ul);
+ /*
+ * We have to do this via a function because the event listener
+ * for the global click needs to be removable, but the click
+ * handler for the menu items also wants to use it.
+ */
+ const closeMenu = () => {
+ if (document.body.contains(menu)) {
+ document.body.removeChild(menu);
+ }
+ window.removeEventListener("click", closeMenu, true);
+ };
+
menuItems.forEach(item => {
const li = document.createElement("li");
li.innerHTML = item.label;
if (item.callback) {
- li.addEventListener("click", item.callback);
+ li.addEventListener("click", (e) => {
+ item.callback(e);
+ closeMenu();
+ });
+ } else {
+ li.classList.add("read-only");
}
ul.appendChild(li);
});
document.body.appendChild(menu);
- const closeMenu = () => {
- if (document.body.contains(menu)) {
- document.body.removeChild(menu);
- }
- document.removeEventListener("click", closeMenu);
- };
-
+ /*
+ * Now we have the content, we can see how big it is. If it
+ * is going off the right of the page, move it left so it ends
+ * at the click coordinates.
+ */
+
+ const rect = menu.getBoundingClientRect();
+ if (rect.right > window.innerWidth)
+ menu.style.left = (event.pageX - rect.width) + "px";
+
+ /*
+ * defer adding the click listener so the current click
+ * doesn't trigger it. Use capture on window so we get
+ * it even if the click target stops propagation.
+ */
setTimeout(() => {
- document.addEventListener("click", closeMenu);
+ window.addEventListener("click", closeMenu, true);
}, 0);
}
@@ -1233,20 +1271,22 @@ function createBuilderDiv(plat) {
let plat_os = plat_parts[0] || 'generic';
let plat_arch = plat_parts[1] || 'generic';
let plat_tc = plat_parts[2] || 'generic';
+ let short_name = plat.name.split('.')[0];
let innerHTML = `<table class="nomar"><tbody><tr><td class="bn">`;
- innerHTML += `<img class="ip1 zup" data-sai-src="/sai/${plat_os}.svg">`;
- innerHTML += `<img class="ip1 tread1" data-sai-src="/sai/arch-${plat_arch}.svg">`;
- innerHTML += `<img class="ip1 tread2" data-sai-src="/sai/tc-${plat_tc}.svg">`;
+ innerHTML += `<div class="builder-name-row">` +
+ `<div class="builder-short-name">${hsanitize(short_name)}</div>` +
+ `<div class="builder-icons">` +
+ `<img class="ip1 zup" data-sai-src="/sai/${plat_os}.svg">` +
+ `<img class="ip1 tread1" data-sai-src="/sai/arch-${plat_arch}.svg">` +
+ `<img class="ip1 tread2" data-sai-src="/sai/tc-${plat_tc}.svg">` +
+ `</div></div>`;
innerHTML += `<div class="resource-bars">` +
`<div class="res-bar"><div class="res-bar-inner res-bar-cpu w-0"></div></div>` +
`<div class="res-bar"><div class="res-bar-inner res-bar-ram w-0"></div></div>` +
`<div class="res-bar"><div class="res-bar-inner res-bar-disk w-0"></div></div>` +
`</div>`;
innerHTML += `${plat.peer_ip}` + " " + plat.stay_on;
-// `<div class="server-state">` +
-// `Slots: ${plat.s_avail_slots}, In-flight: ${plat.s_inflight_count}<br>` +
-// `Last Reject: ${plat.s_last_rej_task_uuid ? plat.s_last_rej_task_uuid.substring(0, 8) : 'none'}` + "</div>" +
innerHTML += `</td></tr></tbody></table>`;
platDiv.innerHTML = innerHTML;
@@ -1415,7 +1455,13 @@ function createPconDiv(pcon) {
stats.style.marginLeft = "10px";
stats.style.fontSize = "0.9em";
stats.style.color = "#666";
- stats.textContent = `${d.voltage_v}V ${d.active_power_w}W ${d.current_ma}mA today:${(d.energy_today_wh/1000).toFixed(3)}kWh`;
+ if (d.voltage_v < 70)
+ stats.textContent = "unpowered";
+ else if (!d.active_power_w)
+ stats.textContent = "OFF";
+ else
+ stats.textContent = `${d.active_power_w}W`;
+
header.appendChild(stats);
}
@@ -1466,9 +1512,33 @@ function createPconDiv(pcon) {
return pconDiv;
}
+let last_renderPconHierarchy_state = "";
+
function renderPconHierarchy(container) {
if (!container) return;
+ const builders_no_time = last_builder_list.map(b => {
+ const { last_seen, ...rest } = b;
+ return rest;
+ });
+
+ const clean_pcons = {};
+ for (const [k, p] of Object.entries(pcon_topology)) {
+ const { children, ...rest } = p;
+ clean_pcons[k] = rest;
+ }
+
+ /* Serialize the inputs to quickly see if we actually need to redraw everything */
+ const currentState = JSON.stringify({
+ pcons: clean_pcons,
+ builders: builders_no_time
+ });
+
+ if (currentState === last_renderPconHierarchy_state) {
+ return;
+ }
+ last_renderPconHierarchy_state = currentState;
+
/* Clear and redraw for now to ensure structure is correct */
container.innerHTML = "";
@@ -2019,6 +2089,7 @@ function ws_open_sai()
if (document.getElementById("sai_overview")) {
document.getElementById("sai_overview").innerHTML = s;
+ logs_pending = times_pending = lines_pending = "";
if (document.getElementById("esr-" + jso.e.uuid))
document.getElementById("esr-" + jso.e.uuid).innerHTML =
@@ -2058,6 +2129,7 @@ function ws_open_sai()
document.getElementById("dlogst").innerHTML = "";
document.getElementById("logs").innerHTML = "";
lines = times = logs = "";
+ lines_pending = times_pending = logs_pending = "";
logAnsiState = {};
tfirst = 0;
lli = 1;
@@ -2201,23 +2273,31 @@ function ws_open_sai()
switch (jso.channel) {
case 1:
- logs += s;
+ logs += s; logs_pending += s;
break;
case 2:
logs += "<span class=\"stderr\">" + s +
"</span>";
+ logs_pending += "<span class=\"stderr\">" + s +
+ "</span>";
break;
case 3:
logs += "<span class=\"saibuild\">\u{25a0} " + s +
"</span>";
+ logs_pending += "<span class=\"saibuild\">\u{25a0} " + s +
+ "</span>";
break;
case 4:
logs += "<span class=\"tty0\">" + s +
"</span>";
+ logs_pending += "<span class=\"tty0\">" + s +
+ "</span>";
break;
default:
logs += "<span class=\"tty1\">" + s +
"</span>";
+ logs_pending += "<span class=\"tty1\">" + s +
+ "</span>";
}
@@ -2236,8 +2316,8 @@ function ws_open_sai()
lli++;
}
- lines += en;
- times += tn;
+ lines += en; lines_pending += en;
+ times += tn; times_pending += tn;
if (!redpend) {
redpend = 1;
@@ -2250,13 +2330,20 @@ function ws_open_sai()
rightPane.scrollTop + 1;
if (document.getElementById("logs")) {
- document.getElementById("logs").innerHTML = logs;
+ if (logs_pending) {
+ document.getElementById("logs").insertAdjacentHTML('beforeend', logs_pending);
+ logs_pending = "";
+ }
- if (document.getElementById("dlogsn"))
- document.getElementById("dlogsn").innerHTML = lines;
+ if (document.getElementById("dlogsn") && lines_pending) {
+ document.getElementById("dlogsn").insertAdjacentHTML('beforeend', lines_pending);
+ lines_pending = "";
+ }
- if (document.getElementById("dlogst"))
- document.getElementById("dlogst").innerHTML = times;
+ if (document.getElementById("dlogst") && times_pending) {
+ document.getElementById("dlogst").insertAdjacentHTML('beforeend', times_pending);
+ times_pending = "";
+ }
}
if (locked && rightPane)
@@ -2425,6 +2512,32 @@ window.addEventListener("load", function() {
}
];
+ const isFinalState = ["taskstate3", "taskstate4", "taskstate5", "taskstate7"].some(s => taskDiv.classList.contains(s));
+
+ if (!isFinalState) {
+ if (taskDiv.classList.contains("taskstate10")) {
+ menuItems.push({
+ label: "Continue task",
+ callback: () => {
+ sai.send(JSON.stringify({
+ schema: "com.warmcat.sai.taskresume",
+ uuid: taskUuid
+ }));
+ }
+ });
+ } else {
+ menuItems.push({
+ label: "Pause task",
+ callback: () => {
+ sai.send(JSON.stringify({
+ schema: "com.warmcat.sai.taskpause",
+ uuid: taskUuid
+ }));
+ }
+ });
+ }
+ }
+
if (taskDiv.dataset.rebuildable === "1")
menuItems.splice(1, 0, {
label: "Rebuild last step",
diff --git a/src/builder/b-deletion.c b/src/builder/b-deletion.c
index dd821d9..b706127 100644
--- a/src/builder/b-deletion.c
+++ b/src/builder/b-deletion.c
@@ -105,21 +105,31 @@ sai_deletion_worker(const char *home_dir)
if (!p)
continue;
+ lwsl_info("%s: received delete request for '%s'\n", __func__, line);
+
{
struct lws_dir_info di;
char full_path[PATH_MAX];
+ struct stat st;
- memset(&di, 0, sizeof(di));
lws_snprintf(full_path, sizeof(full_path),
"%s/jobs/%s", home_dir, line);
+ if (stat(full_path, &st)) {
+ // lwsl_notice("%s: %s already gone or inaccessible\n", __func__, full_path);
+ continue;
+ }
+
+ memset(&di, 0, sizeof(di));
di.dirpath = full_path;
di.cb = lws_dir_rm_rf_cb;
di.do_toplevel_cb = 1;
+ lwsl_info("%s: performing rm -rf %s\n", __func__, full_path);
+
if (lws_dir_via_info(&di))
- lwsl_err("%s: failed to delete %s\n",
- __func__, full_path);
+ lwsl_err("%s: failed to delete %s: %s\n",
+ __func__, full_path, strerror(errno));
}
} while (1);
@@ -154,21 +164,30 @@ scan_jobs_dir_cb(const char *dirpath, void *user, struct lws_dir_entry *lde)
char path[512];
struct stat sb;
- if (lde->type != LDOT_DIR || lde->name[0] == '.')
+ if (lde->name[0] == '.')
return 0;
lws_start_foreach_dll(struct lws_dll2 *, p, active->owner.head) {
struct active_job_uuid *aj = lws_container_of(p, struct active_job_uuid, list);
- if (!strcmp(aj->uuid, lde->name))
+ if (!strcmp(aj->uuid, lde->name)) {
/* it's an active job, leave it alone */
+ lwsl_info("%s: %s is active\n", __func__, lde->name);
return 0;
+ }
} lws_end_foreach_dll(p);
lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name);
- if (stat(path, &sb))
+ if (stat(path, &sb)) {
+ lwsl_notice("%s: stat failed %s\n", __func__, path);
+ return 0;
+ }
+
+ if (!S_ISDIR(sb.st_mode)) {
+ lwsl_notice("%s: %s is not a dir\n", __func__, path);
return 0;
+ }
/* older than 24h? */
@@ -179,8 +198,9 @@ scan_jobs_dir_cb(const char *dirpath, void *user, struct lws_dir_entry *lde)
DWORD written;
#endif
- lwsl_notice("%s: requesting removal of old job dir %s\n",
- __func__, path);
+ lwsl_info("%s: requesting removal of old job dir %s (age %llus)\n",
+ __func__, path, (unsigned long long)
+ ((uint64_t)lws_now_secs() - (uint64_t)sb.st_mtime));
#if !defined(WIN32)
if (write(builder.pipe_master_wr, temp, LWS_POSIX_LENGTH_CAST(len)) != (ssize_t)len)
@@ -190,6 +210,10 @@ scan_jobs_dir_cb(const char *dirpath, void *user, struct lws_dir_entry *lde)
#endif
lwsl_err("%s: failed to write to deletion worker\n",
__func__);
+ } else {
+ lwsl_info("%s: %s is only %llus old\n", __func__, path,
+ (unsigned long long)
+ ((uint64_t)lws_now_secs() - (uint64_t)sb.st_mtime));
}
return 0;
@@ -204,6 +228,8 @@ sul_cleanup_jobs_cb(lws_sorted_usec_list_t *sul)
struct lwsac *ac = NULL;
char path[256];
+ lwsl_info("%s: starting periodic cleanup\n", __func__);
+
memset(&active, 0, sizeof(active));
/*
diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c
index 816a861..b168fa1 100644
--- a/src/builder/b-nspawn.c
+++ b/src/builder/b-nspawn.c
@@ -24,7 +24,6 @@
#include <libwebsockets.h>
#include <string.h>
-#include <signal.h>
#include <sys/stat.h>
#include <sys/types.h>
@@ -172,6 +171,16 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si,
char s[256];
int n;
+ if (!ns) {
+ if (op) {
+ lwsl_warn("%s: op %p has no ns (orphaned), freeing op\n", __func__, op);
+ if (op->spawn)
+ free(op->spawn);
+ free(op);
+ }
+ return;
+ }
+
saib_log_chunk_create(ns, ">saib> <=== Reaping build process\n", 34, 3);
#if !defined(WIN32)
diff --git a/src/builder/b-power.c b/src/builder/b-power.c
index fc05ec3..37f14d1 100644
--- a/src/builder/b-power.c
+++ b/src/builder/b-power.c
@@ -252,12 +252,12 @@ saib_stay_init(void)
if (lws_ss_create(builder.context, 0, &ssi_saib_power_stay_t,
NULL, &builder.ss_stay, NULL, NULL)) {
- lwsl_err("%s: failed to create sai-power-stay ss\n", __func__);
- return 1;
+ lwsl_err("%s: failed to create sai-power-stay ss (ignoring)\n", __func__);
+ return 0;
}
if (!builder.url_sai_power)
- return 1;
+ return 0;
if (!suspender_exists)
return LWSSSSRET_OK;
diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c
index 3f968b8..40e49f3 100644
--- a/src/builder/b-sai.c
+++ b/src/builder/b-sai.c
@@ -323,16 +323,51 @@ app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link,
*/
switch (target) {
+ case LWS_SYSTATE_CONTEXT_CREATED:
+ {
+ struct lws_context_creation_info info;
+
+ builder.context = mgr->context;
+
+ /*
+ * We have the context, but we haven't dropped privs yet.
+ *
+ * We need to init the builder vhost, which has the pipes for
+ * the suspender and the metrics client on it, and the
+ * suspender itself.
+ */
+
+ memset(&info, 0, sizeof(info));
+ pvo1a.value = builder.metrics_uri;
+ pvo1b.value = builder.metrics_path;
+ pvo1c.value = builder.metrics_secret;
+ info.pvo = &pvo1;
+ info.pprotocols = pprotocols;
+
+ builder.vhost = lws_create_vhost(builder.context, &info);
+ if (!builder.vhost) {
+ lwsl_err("Failed to create tls vhost\n");
+ return 1;
+ }
+
+ saib_power_init();
+
+#if defined(__linux__) || defined(__NetBSD__) || defined(__APPLE__)
+ if (saib_suspender_fork(argv0))
+ return 1;
+#endif
+ break;
+ }
+
case LWS_SYSTATE_OPERATIONAL:
if (current != LWS_SYSTATE_OPERATIONAL)
break;
if (saib_deletion_init(argv0))
return 1;
-#if defined(__APPLE__)
- if (saib_suspender_fork(argv0))
+ if (saib_deletion_init(argv0))
return 1;
-#endif
+
/*
* The builder JSON conf listed servers we want to connect to,
@@ -377,11 +412,13 @@ app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link,
} lws_end_foreach_dll(pxx);
+ lwsl_info("%s: platform config completed, calling saib_stay_init\n", __func__);
if (saib_stay_init())
return 1;
+ lwsl_info("%s: scheduling initial cleanup in 100ms\n", __func__);
lws_sul_schedule(builder.context, 0, &builder.sul_cleanup_jobs,
- sul_cleanup_jobs_cb, SAI_CLEANUP_JOBS_INTERVAL_US);
+ sul_cleanup_jobs_cb, 100 * LWS_US_PER_MS);
/* let's sample the best possible free RAM + disk situation,
* we will derate it a bit when using it */
@@ -566,10 +603,13 @@ saib_app_run(int argc, const char **argv)
info.port = CONTEXT_PORT_NO_LISTEN;
info.pprotocols = pprotocols;
+ info.pprotocols = pprotocols;
+
info.uid = sb.st_uid;
info.gid = sb.st_gid;
+
#if !defined(LWS_WITHOUT_EXTENSIONS)
if (!lws_cmdline_option(argc, argv, "-n"))
info.extensions = extensions;
@@ -601,37 +641,20 @@ saib_app_run(int argc, const char **argv)
/* ... and our vhost... */
- pvo1a.value = builder.metrics_uri;
- pvo1b.value = builder.metrics_path;
- pvo1c.value = builder.metrics_secret;
- info.pvo = &pvo1;
-
- builder.vhost = lws_create_vhost(builder.context, &info);
- if (!builder.vhost) {
- lwsl_err("Failed to create tls vhost\n");
- goto bail;
- }
-
- saib_power_init();
-
-#if defined(__linux__)
- if (//builder.power_off_type &&
- //!strcmp(builder.power_off_type, "suspend") &&
- saib_suspender_fork(argv[0]))
+ builder.context = lws_create_context(&info);
+ if (!builder.context) {
+ lwsl_err("lws init failed\n");
return 1;
-#endif
+ }
-#if defined(__NetBSD__)
- if (saib_suspender_fork(argv[0]))
- return 1;
-#endif
+ /* ... and our vhost... */
while (!lws_service(builder.context, 0) && !interrupted)
;
-bail:
suspender_destroy();
+
/* destroy the unique servers */
lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1,
diff --git a/src/builder/b-task.c b/src/builder/b-task.c
index ea575ea..181d1a0 100644
--- a/src/builder/b-task.c
+++ b/src/builder/b-task.c
@@ -20,8 +20,6 @@
*/
#include <libwebsockets.h>
-#include <string.h>
-#include <signal.h>
#include <sys/stat.h>
#include <assert.h>
#include <fcntl.h>
@@ -409,12 +407,36 @@ saib_sub_cleaner_cb(lws_sorted_usec_list_t *sul)
if (ns->op && ns->op->lsp) {
- lwsl_notice("%s: +++++++++++ killing child process\n", __func__);
+ lwsl_notice("%s: +++++++++++ killing child process (budget %d)\n", __func__, ns->term_budget);
lws_spawn_piped_kill_child_process(ns->op->lsp);
- } else {
- lwsl_err("%s: ============= unable to kill child process -> destroying ns\n", __func__);
- saib_task_destroy(ns);
+
+ if (!ns->term_budget)
+ ns->term_budget = 10;
+
+ /* give it a few goes to react to the signal */
+ if (--ns->term_budget) {
+ lws_sul_schedule(builder.context, 0, &ns->sul_cleaner,
+ saib_sub_cleaner_cb, 250 * LWS_US_PER_MS);
+ return;
+ }
+
+ lwsl_err("%s: ============= unable to kill child process -> destroying ns forcibly\n", __func__);
+ /*
+ * It refused to die after a few seconds... we are giving up on it.
+ * Break the link between the op and the ns, so if the op and its
+ * process ever do die, the reap callback will see ns is NULL and
+ * just free the op.
+ */
+ ns->op->ns = NULL;
+ /*
+ * And lose our link to the op, so saib_task_destroy() doesn't
+ * try to kill it again. Because op->ns is NULL, we are leaving
+ * responsibility for freeing op to the eventual reap action.
+ */
+ ns->op = NULL;
}
+
+ saib_task_destroy(ns);
}
void
@@ -720,7 +742,8 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a,
struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list);
if (xns->task && !strcmp(xns->task->uuid, task->uuid)) {
- lwsl_warn("%s: server offered task that's already running\n", __func__);
+ lwsl_warn("%s: server offered task that's already running. State %d, artifacts %d, op %p\n",
+ __func__, xns->state, xns->count_artifacts, xns->op);
saib_queue_task_status_update(sp, spm, task->uuid, 0,
SAI_TASK_REASON_DUPE);
saib_reassess_idle_situation();
diff --git a/src/common/c-utils.c b/src/common/c-utils.c
index 46c1fdf..473d798 100644
--- a/src/common/c-utils.c
+++ b/src/common/c-utils.c
@@ -142,7 +142,7 @@ sai_ss_serialize_queue_helper(struct lws_ss_handle *h,
r = lws_struct_json_serialize(js, buf + LWS_PRE,
sizeof(buf) - LWS_PRE, &w);
- lwsl_hexdump_err(buf + LWS_PRE, w);
+ // lwsl_hexdump_err(buf + LWS_PRE, w);
sai_ss_queue_frag_on_buflist_REQUIRES_LWS_PRE(h, buflist,
buf + LWS_PRE, w, (unsigned int)((fi ? LWSSS_FLAG_SOM : 0) |
@@ -178,6 +178,9 @@ sai_ss_tx_from_buflist_helper(struct lws_ss_handle *ss, struct lws_buflist **buf
if (!(depi & LWSSS_FLAG_SOM))
som = 0;
+ if (*len > fsl)
+ *len = fsl;
+
used = (size_t)lws_buflist_fragment_use(buflist, (uint8_t *)buf, *len, &som1, &eom);
if (!used)
return LWSSSSRET_TX_DONT_SEND;
diff --git a/src/common/include/private.h b/src/common/include/private.h
index d887338..8e96344 100644
--- a/src/common/include/private.h
+++ b/src/common/include/private.h
@@ -21,6 +21,9 @@
* structs common across the various different sai daemons and tools
*/
+#include <libwebsockets.h>
+#include <sqlite3.h>
+
#if defined(WIN32)
#define HAVE_STRUCT_TIMESPEC
#endif
@@ -53,6 +56,7 @@ typedef enum {
SAIES_NOT_READY_FOR_BUILD = 8,
SAIES_STEP_SUCCESS = 9,
+ SAIES_PAUSED = 10,
} sai_event_state_t;
enum {
@@ -192,7 +196,9 @@ struct sai_nspawn {
struct saib_opaque_spawn *op;
sai_task_t *task;
+#if defined(LWS_WITH_SPAWN)
lws_spawn_resource_us_t res;
+#endif
lws_sorted_usec_list_t sul_cleaner;
lws_sorted_usec_list_t sul_mirror;
@@ -597,6 +603,7 @@ typedef struct sai_build_metric {
typedef struct sai_stay {
lws_dll2_t list;
char builder_name[64];
+ char pcon_name[64];
char stay_on; /* 0 = release, 1 = set */
} sai_stay_t;
@@ -687,7 +694,7 @@ typedef struct sai_pcon_control {
*/
extern const lws_struct_map_t
- lsm_stay[2],
+ lsm_stay[3],
lsm_schema_stay[1],
lsm_power_managed_builder[2],
lsm_power_managed_builders_list[2],
diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c
index 8f3eed1..845c29c 100644
--- a/src/common/struct-metadata.c
+++ b/src/common/struct-metadata.c
@@ -302,6 +302,7 @@ const lws_struct_map_t lsm_schema_sq3_map_artifact[] = {
const lws_struct_map_t lsm_stay[] = {
LSM_CARRAY(sai_stay_t, builder_name, "builder_name"),
+ LSM_CARRAY(sai_stay_t, pcon_name, "pcon_name"),
LSM_UNSIGNED(sai_stay_t, stay_on, "stay_on"),
};
diff --git a/src/power/p-http-api.c b/src/power/p-http-api.c
index 5dc6a47..1c8815c 100644
--- a/src/power/p-http-api.c
+++ b/src/power/p-http-api.c
@@ -39,7 +39,9 @@
#include "p-private.h"
+#if defined(LWS_WITH_SPAWN)
extern struct lws_spawn_piped *lsp_wol;
+#endif
extern struct sai_power power;
@@ -71,7 +73,7 @@ find_pcon_by_builder_name(struct sai_power *pwr, const char *builder_name)
lws_start_foreach_dll(struct lws_dll2 *, b_node, pc->registered_builders_owner.head) {
saip_builder_t *sb = lws_container_of(b_node, saip_builder_t, list);
- lwsl_notice("%s: %s %s\n", __func__, sb->name, builder_name);
+// lwsl_notice("%s: %s %s\n", __func__, sb->name, builder_name);
if (!strcmp(sb->name, builder_name))
return pc;
@@ -97,8 +99,10 @@ saip_builder_bringup(saip_server_t *sps, saip_pcon_t *pc)
if (pc->type && !strcmp(pc->type, "wol")) {
if (pc->mac) {
lwsl_notice("%s: triggering WOL for %s\n", __func__, pc->name);
+#if defined(LWS_WITH_SPAWN)
write(lws_spawn_get_fd_stdxxx(lsp_wol, 0),
pc->mac, strlen(pc->mac));
+#endif
} else {
lwsl_err("%s: WOL type but no MAC for %s\n", __func__, pc->name);
}
@@ -149,6 +153,7 @@ saip_set_stay(const char *builder_name, int stay_on)
saip_pcon_t *pc = find_pcon_by_builder_name(&power, builder_name);
/* saip_server_link_t *pss; */ /* Unused? */
saip_server_t *sps;
+ int effective;
if (!pc) {
lwsl_warn("%s: Unknown builder %s\n", __func__, builder_name);
@@ -158,13 +163,19 @@ saip_set_stay(const char *builder_name, int stay_on)
sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list);
/* pss = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); */
- pc->user_keep_on = (char)stay_on;
- saip_notify_server_stay_state(builder_name, stay_on | pc->needed);
+ if (stay_on)
+ pc->flags |= SAIP_PCON_F_MANUAL_STAY;
+ else
+ pc->flags &= (uint8_t)~SAIP_PCON_F_MANUAL_STAY;
+
+ effective = !!(pc->flags & (SAIP_PCON_F_MANUAL_STAY | SAIP_PCON_F_NEEDED));
+
+ saip_notify_server_stay_state(builder_name, effective);
/* Trigger state re-eval */
saip_pcon_start_check();
- if (stay_on | pc->needed) {
+ if (effective) {
/* Ensure it's on immediately if needed */
saip_builder_bringup(sps, pc);
} else {
@@ -405,7 +416,7 @@ local_srv_state(void *userobj, void *sh, lws_ss_constate_t state,
if (pc)
g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload),
- "%c", '0' + (pc->user_keep_on | pc->needed));
+ "%c", '0' + !!(pc->flags & (SAIP_PCON_F_MANUAL_STAY | SAIP_PCON_F_NEEDED)));
else
g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload),
"unknown builder %s", pn);
@@ -429,6 +440,7 @@ local_srv_state(void *userobj, void *sh, lws_ss_constate_t state,
saip_notify_server_power_state(pc->name, 1, 0);
if (pc->mac) {
+#if defined(LWS_WITH_SPAWN)
if (write(lws_spawn_get_fd_stdxxx(lsp_wol, 0),
pc->mac, strlen(pc->mac)) !=
(ssize_t)strlen(pc->mac))
@@ -437,7 +449,8 @@ local_srv_state(void *userobj, void *sh, lws_ss_constate_t state,
else
g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload),
"Resumed %s with stay", pn);
- pc->user_keep_on = 1;
+#endif
+ pc->flags |= SAIP_PCON_F_MANUAL_STAY;
goto bail;
}
@@ -456,7 +469,7 @@ local_srv_state(void *userobj, void *sh, lws_ss_constate_t state,
lwsl_warn("%s: powered on %s\n", __func__, pc->name);
- pc->user_keep_on = 1; /* so builder can understand it's manual */
+ pc->flags |= SAIP_PCON_F_MANUAL_STAY; /* so builder can understand it's manual */
saip_notify_server_power_state(pc->name, 1, 0);
sps = lws_container_of(power.sai_server_owner.head,
@@ -528,11 +541,11 @@ power_off:
} lws_end_foreach_dll(px1);
*/
- if (needs[0] || pc->needed) {
+ if (needs[0] || (pc->flags & SAIP_PCON_F_NEEDED)) {
g->size = (size_t)lws_snprintf(g->payload,
sizeof(g->payload),
"NAK: %s needed: %d, deps needed: '%s'",
- pn, pc->needed, needs);
+ pn, !!(pc->flags & SAIP_PCON_F_NEEDED), needs);
goto bail;
}
}
@@ -554,7 +567,7 @@ power_off:
pc->name,
(int)(SAI_POWERDOWN_HOLDOFF_US / LWS_USEC_PER_SEC));
- pc->user_keep_on = 0; /* reset any manual power up */
+ pc->flags &= (uint8_t)~SAIP_PCON_F_MANUAL_STAY; /* reset any manual power up */
}
bail:
diff --git a/src/power/p-private.h b/src/power/p-private.h
index 83ae36a..0c56d0d 100644
--- a/src/power/p-private.h
+++ b/src/power/p-private.h
@@ -84,9 +84,11 @@ typedef struct saip_pcon {
size_t monitor_rx_pos;
char on;
- char user_keep_on; /* user asked to keep this PCON on via UI */
- char server_requested_on; /* server requested stay for jobs */
- char needed; /* transiently set by deps analysis */
+
+#define SAIP_PCON_F_MANUAL_STAY (1 << 0)
+#define SAIP_PCON_F_NEEDED (1 << 1)
+
+ uint8_t flags;
} saip_pcon_t;
/* Represents a builder connected to us */
diff --git a/src/power/p-sai.c b/src/power/p-sai.c
index 1bb887b..a3b19af 100644
--- a/src/power/p-sai.c
+++ b/src/power/p-sai.c
@@ -77,7 +77,9 @@ int getpid(void) { return 0; }
static const char *config_dir = "/etc/sai/power";
static int interrupted;
static lws_state_notify_link_t nl;
+#if defined(LWS_WITH_SPAWN)
struct lws_spawn_piped *lsp_wol;
+#endif
struct sai_power power;
@@ -211,18 +213,18 @@ sul_pcon_check_cb(lws_sorted_usec_list_t *sul)
}
/* Rule 2: User Keep On -> Turn ON */
- if (pc->user_keep_on) {
+ if (pc->flags & SAIP_PCON_F_MANUAL_STAY) {
target_on = 1;
lwsl_warn("%s: PCON %s has user keep on -> Force ON\n", __func__, pc->name);
}
/* Rule 3: Server Requested -> Turn ON */
- else if (pc->server_requested_on) {
+ else if (pc->flags & SAIP_PCON_F_NEEDED) {
target_on = 1;
lwsl_warn("%s: PCON %s has server request -> Force ON\n", __func__, pc->name);
}
- lwsl_info("%s: PCON %s check: target=%d, current=%d (user=%d, srv=%d)\n",
- __func__, pc->name, target_on, pc->on, pc->user_keep_on, pc->server_requested_on);
+ lwsl_info("%s: PCON %s check: target=%d, current=%d (flags=0x%x)\n",
+ __func__, pc->name, target_on, pc->on, pc->flags);
/* If we decide it should be ON, trigger it */
if (target_on && !pc->on) {
@@ -271,10 +273,10 @@ sul_broadcast_energy_cb(lws_sorted_usec_list_t *sul)
/* SS request triggers the HTTP GET */
if (lws_ss_request_tx(pc->ss_tasmota_monitor))
lwsl_warn("%s: Failed to trigger monitor request for %s\n", __func__, pc->name);
- else
- lwsl_notice("%s: Triggered polling for %s\n", __func__, pc->name);
- } else {
- lwsl_warn("%s: PCON %s has no monitor SS\n", __func__, pc->name);
+ // else
+ // lwsl_notice("%s: Triggered polling for %s\n", __func__, pc->name);
+ //} else {
+ // lwsl_warn("%s: PCON %s has no monitor SS\n", __func__, pc->name);
}
} lws_end_foreach_dll(p);
@@ -289,7 +291,7 @@ sul_broadcast_energy_cb(lws_sorted_usec_list_t *sul)
saip_server_t *sps = lws_container_of(mp, struct saip_server, list);
int queued = saip_queue_energy_report(sps);
if (queued) {
- lwsl_notice("%s: Queued energy report for server\n", __func__);
+ // lwsl_notice("%s: Queued energy report for server\n", __func__);
if (lws_ss_request_tx(sps->ss)) /* Request write to send the report */
lwsl_warn("%s: failed to request tx\n", __func__);
}
@@ -532,6 +534,7 @@ int main(int argc, const char **argv)
saip_ss_create_tasmota();
+#if defined(LWS_WITH_SPAWN)
{
struct lws_spawn_piped_info info;
char rpath[PATH_MAX];
@@ -550,8 +553,9 @@ int main(int argc, const char **argv)
if (!lsp_wol)
lwsl_err("%s: wol spawn failed\n", __func__);
}
+#endif
- lws_finalize_startup(power.context);
+ lws_finalize_startup(power.context, "sai-power");
while (!lws_service(power.context, 0) && !interrupted)
diff --git a/src/power/p-smartplug.c b/src/power/p-smartplug.c
index 58a496e..f2265aa 100644
--- a/src/power/p-smartplug.c
+++ b/src/power/p-smartplug.c
@@ -67,8 +67,8 @@ saip_spc_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
/* Success, update latest data and timestamp */
pc->latest_data = tp.td;
pc->last_monitor_time = lws_now_usecs();
- lwsl_notice("%s: Parsed Tasmota data for %s: %u W, %u V\n",
- __func__, pc->name, pc->latest_data.active_power_w, pc->latest_data.voltage_v);
+ // lwsl_notice("%s: Parsed Tasmota data for %s: %u W, %u V\n",
+ // __func__, pc->name, pc->latest_data.active_power_w, pc->latest_data.voltage_v);
} else
lwsl_warn("%s: Failed to parse Tasmota data for %s (ret %d)\n", __func__, pc->name, parse_ret);
diff --git a/src/power/p-utils.c b/src/power/p-utils.c
index c8741be..713faff 100644
--- a/src/power/p-utils.c
+++ b/src/power/p-utils.c
@@ -63,10 +63,40 @@ void
saip_switch(saip_pcon_t *pc, int on)
{
struct lws_ss_handle *h = on ? pc->ss_tasmota_on : pc->ss_tasmota_off;
+ int wol_fired = 0;
+
+ if (on && pc->mac) {
+ char buf[64];
+ size_t n = (size_t)lws_snprintf(buf, sizeof(buf), "%s\n", pc->mac);
+
+ if (write(lws_spawn_get_fd_stdxxx(lsp_wol, 0),
+ buf, n) != (ssize_t)n)
+ lwsl_err("%s: Write to resume %s failed %d\n",
+ __func__, pc->name, errno);
+ else {
+ lwsl_notice("%s: Resumed %s via WOL\n",
+ __func__, pc->name);
+ wol_fired = 1;
+ }
+ }
if (!h) {
- lwsl_err("%s: %s: no ss handle for %s\n", __func__,
- pc->name, on ? "ON" : "OFF");
+ if (!wol_fired) {
+ if (pc->type && !strcmp(pc->type, "wol")) {
+ if (on && !pc->mac)
+ lwsl_err("%s: %s: WOL type but no MAC configured\n",
+ __func__, pc->name);
+ else
+ if (!on)
+ lwsl_info("%s: %s: WOL pcon ignoring OFF\n",
+ __func__, pc->name);
+ } else {
+ lwsl_err("%s: %s: no ss handle for %s (type: %s, mac: %s)\n",
+ __func__, pc->name, on ? "ON" : "OFF",
+ pc->type ? pc->type : "null",
+ pc->mac ? pc->mac : "null");
+ }
+ }
return;
}
diff --git a/src/power/p-ws-server.c b/src/power/p-ws-server.c
index c9754d3..6ad2e64 100644
--- a/src/power/p-ws-server.c
+++ b/src/power/p-ws-server.c
@@ -44,10 +44,6 @@ static const lws_struct_map_t lsm_saip_rx_map[] = {
"com.warmcat.sai.pcon_control"),
};
-/*
- * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c)
- */
-
int
saip_queue_energy_report(saip_server_t *sps)
{
@@ -82,13 +78,13 @@ saip_queue_energy_report(saip_server_t *sps)
} else {
if (pc->last_monitor_time)
lwsl_notice("%s: Stale monitor data for %s (age %llus)\n", __func__, pc->name, (unsigned long long)(lws_now_usecs() - pc->last_monitor_time) / LWS_US_PER_SEC);
- else
- lwsl_notice("%s: No monitor data for %s\n", __func__, pc->name);
+ // else
+ // lwsl_notice("%s: No monitor data for %s\n", __func__, pc->name);
}
} lws_end_foreach_dll(p);
if (count) {
- lwsl_notice("%s: Queuing energy report with %d items\n", __func__, count);
+ // lwsl_notice("%s: Queuing energy report with %d items\n", __func__, count);
r = sai_ss_serialize_queue_helper(sps->ss, &m->bl_pwr_to_srv,
lsm_schema_pcon_energy,
LWS_ARRAY_SIZE(lsm_schema_pcon_energy),
@@ -201,9 +197,8 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
lws_struct_args_t a;
struct lejp_ctx ctx;
- lwsl_notice("%s: len %d, flags: %d (saip_server_t %p)\n", __func__, (int)len, flags, (void *)sps);
+ lwsl_notice("%s: PPPPPPPP len %d, flags: %d (saip_server_t %p)\n", __func__, (int)len, flags, (void *)sps);
lwsl_hexdump_notice(buf, len);
- /* lwsl_hexdump_notice(buf, len); */
memset(&a, 0, sizeof(a));
a.map_st[0] = lsm_saip_rx_map;
@@ -221,12 +216,14 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
lwsl_warn("%s: RX PCON Control '%s' -> %d\n", __func__, ctl->pcon_name, ctl->on);
if (pc) {
- lwsl_warn("%s: Applying PCON Control '%s' -> %d (prev user_keep_on=%d)\n",
- __func__, pc->name, ctl->on, pc->user_keep_on);
- pc->user_keep_on = ctl->on;
+ lwsl_warn("%s: Applying PCON Control '%s' -> %d (prev flags=0x%x)\n",
+ __func__, pc->name, ctl->on, pc->flags);
+
if (ctl->on) {
+ pc->flags |= SAIP_PCON_F_MANUAL_STAY;
saip_switch(pc, 1);
} else {
+ pc->flags &= (uint8_t)~SAIP_PCON_F_MANUAL_STAY;
saip_pcon_start_check();
}
} else {
@@ -235,6 +232,7 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
} else {
/* Stay */
sai_stay_t *stay = (sai_stay_t *)a.dest;
+ saip_pcon_t *pc;
// {"schema":"com.warmcat.sai.power.stay","builder_name":"ubuntu_rpi4","stay_on":1}
@@ -247,22 +245,42 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
* We should map this back to the PCON.
*/
+ if (stay->pcon_name[0]) {
+ pc = saip_pcon_by_name(&power, stay->pcon_name);
+ if (pc) {
+ lwsl_notice("%s: Direct map stay for PCON '%s'\n",
+ __func__, pc->name);
+ /* Update PCON stay state */
+
+ if (stay->stay_on) {
+ pc->flags |= SAIP_PCON_F_MANUAL_STAY;
+ /* If stay is set, ensure it is on immediately */
+ saip_switch(pc, 1);
+ } else {
+ pc->flags &= (uint8_t)~SAIP_PCON_F_MANUAL_STAY;
+ /* If stay is cleared, schedule power off check */
+ saip_pcon_start_check();
+ }
+ goto found;
+ }
+ }
+
lws_start_foreach_dll(struct lws_dll2 *, p, power.sai_pcon_owner.head) {
- saip_pcon_t *pc = lws_container_of(p, saip_pcon_t, list);
+ pc = lws_container_of(p, saip_pcon_t, list);
lws_start_foreach_dll(struct lws_dll2 *, b_node, pc->registered_builders_owner.head) {
saip_builder_t *sb = lws_container_of(b_node, saip_builder_t, list);
if (!strcmp(sb->name, stay->builder_name)) {
lwsl_notice("%s: Mapping stay for builder '%s' to PCON '%s'\n",
__func__, sb->name, pc->name);
/* Update PCON stay state */
- pc->server_requested_on = stay->stay_on;
-
- /* If stay is cleared, schedule power off check */
- if (!stay->stay_on)
- saip_pcon_start_check();
- else {
+ if (stay->stay_on) {
+ pc->flags |= SAIP_PCON_F_MANUAL_STAY;
/* If stay is set, ensure it is on immediately */
saip_switch(pc, 1);
+ } else {
+ pc->flags &= (uint8_t)~SAIP_PCON_F_MANUAL_STAY;
+ /* If stay is cleared, schedule power off check */
+ saip_pcon_start_check();
}
goto found;
}
@@ -277,15 +295,86 @@ found:
lwsac_free(&a.ac);
/*
- * The old logic parsed comma-separated platform names to determine needed state.
- * We are moving away from that. Sai-server should explicitly request power state
- * or we should rely on the builder being "needed" implies PCON on.
- * Actually, the requirement said: "You can no longer ask that a platform stays on, instead, you ask that the power controller stays on"
- * So sai-server might send stay requests for PCONs directly if we update the UI.
- * But for now, sai-server logic (s-power.c) sends stay for *builders*.
- * So the mapping logic above is correct for transition.
+ * It wasn't a JSON message... it's the comma-separated list of needed
+ * platforms then
*/
+ lwsl_err("%s: ************* Received comma-separated list of needed platforms\n", __func__);
+ sai_dump_stderr(buf, len);
+
+ lws_start_foreach_dll(struct lws_dll2 *, p, power.sai_pcon_owner.head) {
+ saip_pcon_t *pc = lws_container_of(p, saip_pcon_t, list);
+
+ pc->flags &= (uint8_t)~SAIP_PCON_F_NEEDED;
+ } lws_end_foreach_dll(p);
+
+ if (len) {
+ const char *cp = (const char *)buf;
+ const char *end = cp + len;
+
+ while (cp < end) {
+ const char *comma = memchr(cp, ',', lws_ptr_diff_size_t(end, cp));
+ size_t token_len;
+
+ if (comma)
+ token_len = lws_ptr_diff_size_t(comma, cp);
+ else
+ token_len = lws_ptr_diff_size_t(end, cp);
+
+ if (token_len) {
+ char pcon[64];
+ saip_pcon_t *pc;
+
+ lws_strnncpy(pcon, cp, token_len, sizeof(pcon));
+ pc = saip_pcon_by_name(&power, pcon);
+ if (pc)
+ pc->flags |= SAIP_PCON_F_NEEDED;
+ else
+ lwsl_notice("%s: unknown pcon '%.*s' needed\n",
+ __func__, (int)token_len, cp);
+ }
+
+ cp += token_len;
+ if (cp < end && *cp == ',')
+ cp++;
+ }
+ }
+
+ /*
+ * Propagate needed state up the dependency tree
+ *
+ * If a PCON is needed, and it depends on another PCON, that parent PCON
+ * is also needed.
+ */
+ {
+ int changed;
+
+ do {
+ changed = 0;
+ lws_start_foreach_dll(struct lws_dll2 *, p,
+ power.sai_pcon_owner.head) {
+ saip_pcon_t *pc = lws_container_of(p,
+ saip_pcon_t, list);
+
+ if (pc->flags & SAIP_PCON_F_NEEDED) {
+ /* check if this PCON depends on another */
+ if (pc->depends_on) {
+ saip_pcon_t *parent = saip_pcon_by_name(&power,
+ pc->depends_on);
+ if (parent && !(parent->flags & SAIP_PCON_F_NEEDED)) {
+ parent->flags |= SAIP_PCON_F_NEEDED;
+ changed = 1;
+ lwsl_notice("%s: PCON %s needed by dep %s\n",
+ __func__, parent->name, pc->name);
+ }
+ }
+ }
+ } lws_end_foreach_dll(p);
+ } while (changed);
+ }
+
+ saip_pcon_start_check();
+
return 0;
}
diff --git a/src/server/s-central.c b/src/server/s-central.c
index 12a1846..802a062 100644
--- a/src/server/s-central.c
+++ b/src/server/s-central.c
@@ -25,9 +25,7 @@
*/
#include <libwebsockets.h>
-#include <string.h>
-#include <signal.h>
-#include <time.h>
+
#include "s-private.h"
@@ -203,6 +201,13 @@ sais_central_cb(lws_sorted_usec_list_t *sul)
vhd->last_check_abandoned_tasks = lws_now_usecs();
}
+ /*
+ * Wake up builders if there are new tasks that just became mature enough
+ * (older than 10s)
+ */
+ sais_prune_inflight_list(vhd);
+ sais_platforms_with_tasks_pending(vhd);
+
/* check again in 1s */
lws_sul_schedule(context, 0, &vhd->sul_central, sais_central_cb,
diff --git a/src/server/s-notification.c b/src/server/s-notification.c
index 45a5ab9..41ca405 100644
--- a/src/server/s-notification.c
+++ b/src/server/s-notification.c
@@ -583,6 +583,7 @@ next_plat: ;
pss->sn.t.repo_name = pss->sn.e.repo_name;
pss->sn.t.git_ref = sn->e.ref;
pss->sn.t.git_hash = sn->e.hash;
+ pss->sn.t.parallel = 2;
lws_dll2_clear(&pss->sn.t.list);
lws_dll2_owner_clear(&owner);
diff --git a/src/server/s-power.c b/src/server/s-power.c
index 049efc3..4e681e2 100644
--- a/src/server/s-power.c
+++ b/src/server/s-power.c
@@ -37,10 +37,29 @@
#include "s-private.h"
+#if 0
/*
* (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c)
*/
+static int
+sais_get_pcon_for_builder(struct vhd *vhd, const char *builder_name,
+ char *pcon_buf, size_t len)
+{
+ char q[256], esc[96];
+ int r;
+
+ lws_sql_purify(esc, builder_name, sizeof(esc));
+
+ lws_snprintf(q, sizeof(q),
+ "SELECT pcon_name FROM pcon_builders WHERE builder_name = '%s'",
+ esc);
+
+ r = sqlite3_exec(vhd->server.pdb, q, sql3_get_string_cb, pcon_buf, NULL);
+
+ return r == SQLITE_OK;
+}
+#endif
int
sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf,
size_t bl, unsigned int ss_flags)
@@ -194,7 +213,15 @@ passthru:
* We want to broadcast it to all connected web interfaces (sai-web).
* We have been buffering in pss->power_rx_cache until we identified the schema.
* Now we forward whatever is in the cache (which includes the current 'buf').
+ *
+ * Notice we wait until we have the WHOLE message in the cache before
+ * broadcasting it. Since there is only one websrv channel, we can't
+ * allow interleaved fragments from different sources on it.
*/
+
+ if (!(ss_flags & LWSSS_FLAG_EOM))
+ return 0;
+
{
lws_wsmsg_info_t info;
uint8_t *p, *lin;
@@ -231,122 +258,7 @@ passthru:
info.private_source_idx = SAI_WEBSRV_PB__GENERATED;
info.buf = p + LWS_PRE;
info.len = tlen;
- info.ss_flags = LWSSS_FLAG_SOM; /* We always send what we have as a start */
-
- if (ss_flags & LWSSS_FLAG_EOM)
- info.ss_flags |= LWSSS_FLAG_EOM;
-
- /*
- * If we are continuing (n == LEJP_CONTINUE), we flushed the buffer
- * so subsequent calls will append new data to empty buflist and flush it immediately.
- * However, sais_websrv_broadcast expects SOM/EOM to be correct for the whole message.
- *
- * If we buffered the START of the message, we set SOM.
- * If the incoming chunk was EOM, we set EOM.
- *
- * What if we have intermediate chunks?
- *
- * If we are in passthru, we cleared the cache above.
- *
- * Wait, if we are in passthru state, we shouldn't re-set SOM for every chunk.
- * We need to track if we already sent SOM.
- *
- * But here we only enter `passthru` if `top_schema_index` matches.
- * This happens for the FIRST chunk (once schema matches) AND subsequent chunks.
- *
- * Problem: `top_schema_index` remains 3 for subsequent chunks.
- *
- * So we need to know if we are flushing the FIRST part (SOM) or a later part.
- *
- * `ss_flags & LWSSS_FLAG_SOM` tells us if the CURRENT chunk was the start of the message.
- *
- * If `ss_flags & SOM`, then `info.ss_flags |= SOM`.
- * If `ss_flags & EOM`, then `info.ss_flags |= EOM`.
- *
- * This seems correct because we are effectively delaying the processing.
- * If we buffered chunks 1 and 2, and now processing chunk 2 (which made schema valid),
- * chunk 1 had SOM. `pss->power_rx_cache` contains chunk 1 + chunk 2.
- * So the aggregate buffer DOES start with SOM content.
- *
- * If we are processing chunk 3 (schema already known), we append to cache, then flush.
- * Cache contains just chunk 3.
- * Chunk 3 does NOT have SOM.
- * So we shouldn't set SOM.
- *
- * BUT `ss_flags` belongs to the current `buf` (chunk 3).
- * If `ss_flags` has SOM, then our buffer starts with SOM.
- *
- * So `info.ss_flags = ss_flags` is ALMOST correct, except that we might have accumulated
- * previous chunks which HAD SOM, even if the current chunk doesn't.
- *
- * If `fragment_cache` was non-empty before we appended `buf`, then we are continuing a buffer.
- * Wait, we appended `buf` to `cache` at the top of the function.
- *
- * If `ss_flags` has SOM, then the cache definitely starts with SOM.
- *
- * If `ss_flags` does NOT have SOM, but we have older data in cache?
- * That older data MUST be the start of the message (because we flush on schema detection).
- *
- * Wait, if we flush on schema detection, we flush the START.
- * Subsequent chunks will be appended to EMPTY cache, then flushed.
- *
- * So if cache has data, and we are flushing...
- *
- * Case 1: First chunk(s). Schema found. `ss_flags` might be SOM (if single chunk) or NOT (if 2nd chunk).
- * If 2nd chunk triggers match, `ss_flags` is !SOM. But cache contains Chunk 1 (SOM) + Chunk 2.
- * So we must set SOM if the *cache* contains the start.
- *
- * We can track `pss->power_rx_cache_had_som`.
- * Or we can just rely on `a->top_schema_index == -1` -> we are at start.
- * Once matched, we are flushing the start.
- *
- * Actually, simpler:
- * We only buffer if we DON'T know the schema.
- *
- * If we know the schema (3), we are in passthru mode.
- *
- * If we just transitioned to schema 3 (match occurred in this chunk), we flush everything. This flush INCLUDES the start. So send SOM.
- *
- * If we were ALREADY in schema 3 (subsequent chunks), we just forward `buf`.
- *
- * But wait, my logic "always append to cache" means `buf` is in cache.
- *
- * If I flush cache every time `passthru` is hit:
- * - First time (match): Cache has Start + ... + Current. Flush. Send SOM.
- * - Next time: Cache has Next Chunk. Flush. Send !SOM.
- *
- * How do I know if it's the "First time"?
- * `a->top_schema_index` is persistent in `pss`.
- *
- * Valid point: `lws_struct` parser state persists.
- *
- * Issue: `lejp` doesn't tell me "I just matched schema".
- *
- * But I can check if `pss->power_rx_cache` contains more than `bl`.
- * If `total_len > bl`, then we have buffered data -> We are sending the start -> SOM.
- *
- * Exception: What if `buf` is the FIRST chunk and it matched?
- * `total_len == bl`. But `ss_flags` has SOM.
- *
- * So logic:
- * `info.ss_flags = (ss_flags & LWSSS_FLAG_EOM);`
- * `if (total_len > bl || (ss_flags & LWSSS_FLAG_SOM)) info.ss_flags |= LWSSS_FLAG_SOM;`
- *
- * This handles:
- * - Single chunk (SOM+EOM): len==bl, flags=SOM. Result: SOM+EOM.
- * - Multi chunk, 1st (match): len==bl, flags=SOM. Result: SOM.
- * - Multi chunk, 2nd (match): len > bl, flags=!SOM. Result: SOM. (Correct, as it contains start).
- * - Multi chunk, 3rd (already matched):
- * Wait, if already matched, we still append and flush?
- * If we flush every time, cache is empty between calls.
- * So for 3rd chunk, `total_len == bl`. `flags`=!SOM. Result: !SOM. Correct.
- *
- */
-
- if (tlen > bl || (ss_flags & LWSSS_FLAG_SOM))
- info.ss_flags |= LWSSS_FLAG_SOM;
- else
- info.ss_flags &= (unsigned int)~LWSSS_FLAG_SOM;
+ info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM;
/* Broadcast to all websrv connections (i.e. all sai-web instances) */
sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info);
@@ -375,14 +287,44 @@ bail:
* (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c)
*/
+struct pcon_lookup_ctx {
+ char *p;
+ const char *end;
+ int *n;
+};
+
+static int
+cb_lookup_pcon(void *user, int cols, char **values, char **name)
+{
+ struct pcon_lookup_ctx *ctx = (struct pcon_lookup_ctx *)user;
+ size_t m;
+
+ if (cols < 1 || !values[0])
+ return 0;
+
+ m = strlen(values[0]);
+
+ if (*ctx->n)
+ *ctx->p++ = ',';
+
+ if (lws_ptr_diff_size_t(ctx->end, ctx->p) < m + 2)
+ return 1; /* abort */
+
+ memcpy(ctx->p, values[0], m);
+ ctx->p += m;
+ *ctx->p = '\0';
+ *ctx->n = 1;
+
+ return 0;
+}
+
int
sais_power_tx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl)
{
uint8_t *start = buf + LWS_PRE, *p = start, *end = p + bl - LWS_PRE - 1;
enum lws_write_protocol flags;
- char diff = 0;
size_t w;
- int n;
+ int n, diff = 0;
if (pss->stay_owner.head) {
/*
@@ -453,39 +395,58 @@ sais_power_tx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl)
n = 0;
lws_start_foreach_dll(struct lws_dll2 *, px, vhd->pending_plats.head) {
sais_plat_t *pl = lws_container_of(px, sais_plat_t, list);
- size_t m;
+ struct pcon_lookup_ctx ctx;
+ char q[256], query[256];
+ int r;
+
+ ctx.p = (char *)p;
+ ctx.end = (const char *)end;
+ ctx.n = &n;
+
+ lws_sql_purify(q, pl->plat, sizeof(q));
+ lws_snprintf(query, sizeof(query),
+ "SELECT DISTINCT pcon FROM builders WHERE platform = '%s'",
+ q);
- if (n)
- *p++ = ',';
- m = strlen(pl->plat);
- if (lws_ptr_diff_size_t(end, p) < m + 2)
+ r = sqlite3_exec(vhd->server.pdb, query, cb_lookup_pcon, &ctx, NULL);
+
+ lwsl_notice("%s: platform '%s' -> pcon query '%s': result %d\n",
+ __func__, pl->plat, query, r);
+
+ if (r != SQLITE_OK)
+ lwsl_err("%s: sqlite3 error: %s\n", __func__, sqlite3_errmsg(vhd->server.pdb));
+
+ p = (uint8_t *)ctx.p;
+ if (p >= (uint8_t *)end) /* buffer full */
break;
- memcpy(p, pl->plat, m);
- p += m;
- *p = '\0';
- n = 1;
} lws_end_foreach_dll(px);
+ lwsl_notice("%s: final pcon list: '%.*s'\n", __func__,
+ (int)lws_ptr_diff_size_t(p, start), start);
+
/*
* Don't resend the same status over and over
*/
if (strncmp(pss->last_power_report, (const char *)start, lws_ptr_diff_size_t(p, start) + 1)) {
- diff = 1;
+ lwsl_notice("%s: pending plats changed: '%s' -> '%.*s'\n", __func__,
+ pss->last_power_report, (int)lws_ptr_diff_size_t(p, start), start);
+
memcpy(pss->last_power_report, start, lws_ptr_diff_size_t(p, start) + 1);
- }
+ diff = 1;
+ }
- if (diff /* && start != p */) {
- lwsl_notice("%s: detected jobs for %.*s\n", __func__,
- (int)lws_ptr_diff_size_t(p, start), start);
+ if (diff) {
+ lwsl_notice("%s: ************* detected jobs for %.*s\n", __func__,
+ (int)lws_ptr_diff_size_t(p, start), start);
- if (lws_write(pss->wsi, start, lws_ptr_diff_size_t(p, start),
- LWS_WRITE_TEXT) < 0)
- return -1;
+ if (lws_write(pss->wsi, start, lws_ptr_diff_size_t(p, start),
+ LWS_WRITE_TEXT) < 0)
+ return -1;
- lws_callback_on_writable(pss->wsi);
- }
+ lws_callback_on_writable(pss->wsi);
+ }
return 0;
}
diff --git a/src/server/s-private.h b/src/server/s-private.h
index f9c9178..f3b3f7a 100644
--- a/src/server/s-private.h
+++ b/src/server/s-private.h
@@ -312,6 +312,9 @@ sais_set_task_state(struct vhd *vhd, const char *task_uuid, sai_event_state_t st
uint64_t started, uint64_t duration);
int
+sais_task_pause(struct vhd *vhd, const char *task_uuid);
+
+int
sais_bind_task_to_builder(struct vhd *vhd, const char *builder_name,
const char *builder_uuid, const char *task_uuid);
diff --git a/src/server/s-task-helpers.c b/src/server/s-task-helpers.c
index 5efe89c..f2cfec6 100644
--- a/src/server/s-task-helpers.c
+++ b/src/server/s-task-helpers.c
@@ -215,6 +215,12 @@ sais_set_task_state(struct vhd *vhd, const char *task_uuid,
goto bail;
}
+ if (task_ostate == SAIES_PAUSED &&
+ (state == SAIES_BEING_BUILT || state == SAIES_PASSED_TO_BUILDER ||
+ state == SAIES_STEP_SUCCESS || state == SAIES_FAIL ||
+ state == SAIES_CANCELLED || state == SAIES_SUCCESS))
+ state = SAIES_PAUSED;
+
if (started) {
if (started == 1)
lws_snprintf(esc3, sizeof(esc3), ",started=0");
@@ -364,6 +370,54 @@ bail:
}
int
+sais_task_pause(struct vhd *vhd, const char *task_uuid)
+{
+ char event_uuid[33], esc_uuid[129], q[128];
+ int build_step = -1, state = -1;
+ sqlite3 *pdb = NULL;
+
+ sai_task_uuid_to_event_uuid(event_uuid, task_uuid);
+ if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache,
+ vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) {
+ lwsl_err("%s: unable to open db for event %s\n", __func__, event_uuid);
+ return -1;
+ }
+
+ lws_sql_purify(esc_uuid, task_uuid, sizeof(esc_uuid));
+ lws_snprintf(q, sizeof(q),
+ "select build_step,state from tasks where uuid='%s'",
+ esc_uuid);
+
+ if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &build_step,
+ NULL) != SQLITE_OK)
+ build_step = -1;
+ /* sql3_get_integer_cb will put the last column in there */
+ if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &state,
+ NULL) != SQLITE_OK)
+ state = -1;
+
+ if (state == SAIES_PASSED_TO_BUILDER || state == SAIES_BEING_BUILT) {
+ /*
+ * If it's already building, we need to cancel it on the
+ * builder and decrement the step so it restarts this step
+ * on resume
+ */
+ if (build_step > 0) {
+ build_step--;
+ lws_snprintf(q, sizeof(q),
+ "update tasks set build_step=%d where uuid='%s'",
+ build_step, esc_uuid);
+ sqlite3_exec(pdb, q, NULL, NULL, NULL);
+ }
+ sais_task_stop_on_builders(vhd, task_uuid);
+ }
+
+ sai_event_db_close(&vhd->sqlite3_cache, &pdb);
+
+ return sais_set_task_state(vhd, task_uuid, SAIES_PAUSED, 0, 0);
+}
+
+int
sais_task_cancel(struct vhd *vhd, const char *task_uuid)
{
sai_cancel_t *can;
diff --git a/src/server/s-task.c b/src/server/s-task.c
index 02f5352..2d03759 100644
--- a/src/server/s-task.c
+++ b/src/server/s-task.c
@@ -21,8 +21,7 @@
#include <libwebsockets.h>
#include <string.h>
-#include <signal.h>
-#include <time.h>
+
#include <assert.h>
#include "s-private.h"
@@ -156,7 +155,7 @@ sais_prune_inflight_list(struct vhd *vhd)
lws_start_foreach_dll_safe(struct lws_dll2 *, p1, p2, sp->inflight_owner.head) {
sai_uuid_list_t *u = lws_container_of(p1, sai_uuid_list_t, list);
- if (!u->started && (t - u->us_time_listed) > 5 * 1000 * 1000)
+ if (!u->started && (t - u->us_time_listed) > 3 * 1000 * 1000)
sais_inflight_entry_destroy(u);
} lws_end_foreach_dll_safe(p1, p2);
@@ -504,7 +503,8 @@ sais_platforms_with_tasks_pending(struct vhd *vhd)
* Collect a list of *events* (not tasks) that still have any open tasks
*/
- lws_snprintf(pf, sizeof(pf)," and (state != 3 and state != 5)");
+ lws_snprintf(pf, sizeof(pf)," and (state != 3 and state != 5) and (created < %llu)",
+ (unsigned long long)(lws_now_secs() - 10));
n = lws_struct_sq3_deserialize(vhd->server.pdb, pf, "created desc ",
lsm_schema_sq3_map_event, &o, &ac, 0, 20);
@@ -564,6 +564,17 @@ sais_platforms_with_tasks_pending(struct vhd *vhd)
sais_notify_all_sai_power(vhd);
+ /*
+ * Wake up any builders that are slacking, since there are new tasks
+ */
+ lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) {
+ sai_plat_t *sp = lws_container_of(p, sai_plat_t, sai_plat_list);
+
+ if (!sp->busy)
+ sais_plat_busy(sp, 0);
+
+ } lws_end_foreach_dll(p);
+
lwsac_free(&ac);
return 0;
@@ -749,7 +760,7 @@ nope:
info.private_source_idx = SAI_WEBSRV_PB__ACTIVITY;
info.buf = (uint8_t *)start;
info.len = lws_ptr_diff_size_t(p, start);
- info.ss_flags = LWSSS_FLAG_EOM;
+ info.ss_flags = (unsigned int)((s ? LWSSS_FLAG_SOM : 0) | LWSSS_FLAG_EOM);
sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info);
lws_sul_schedule(vhd->context, 0, &vhd->sul_activity,
diff --git a/src/server/s-webops.c b/src/server/s-webops.c
index 05e87c2..f2f9444 100644
--- a/src/server/s-webops.c
+++ b/src/server/s-webops.c
@@ -22,7 +22,6 @@
#include <libwebsockets.h>
#include <string.h>
-#include <signal.h>
#include <time.h>
#include <assert.h>
@@ -55,16 +54,15 @@ static void
_sais_websrv_broadcast(struct lws_ss_handle *h, void *arg)
{
websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h);
- lws_wsmsg_info_t *info = (lws_wsmsg_info_t *)arg;
- unsigned int *pi = (unsigned int *)((const char *)info->buf - sizeof(int));
+ lws_wsmsg_info_t *info_in = (lws_wsmsg_info_t *)arg;
+ lws_wsmsg_info_t info = *info_in;
+ unsigned int *pi = (unsigned int *)((const char *)info.buf - sizeof(int));
- info->head_upstream = &m->bl_srv_to_web;
- info->private_heads = m->private_heads;
+ info.head_upstream = &m->bl_srv_to_web;
+ info.private_heads = m->private_heads;
// lwsl_ss_notice(h, "Queueing %u bytes, ridx %d, ff_flags: %u",
- // (unsigned int)info->len, info->private_source_idx, info->ss_flags);
-
- *pi = info->ss_flags;
+ // (unsigned int)info.len, info.private_source_idx, info.ss_flags);
/* sai-web might not be taking it.. */
@@ -76,10 +74,12 @@ _sais_websrv_broadcast(struct lws_ss_handle *h, void *arg)
return;
}
- info->buf = info->buf - sizeof(int);
- info->len = info->len + sizeof(int);
+ *pi = info.ss_flags;
+
+ info.buf = info.buf - sizeof(int);
+ info.len = info.len + sizeof(int);
- if (lws_wsmsg_append(info) < 0)
+ if (lws_wsmsg_append(&info) < 0)
lwsl_ss_err(h, "failed to append"); /* still ask to drain */
if (lws_ss_request_tx(h))
@@ -97,21 +97,6 @@ sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv,
}
-struct sai_bl_args {
- uint8_t *buf;
- size_t len;
-};
-
-static void
-_sais_websrv_broadcast_buflist(struct lws_ss_handle *h, void *arg)
-{
- websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h);
- struct sai_bl_args *sbba = (struct sai_bl_args *)arg;
-
- if (lws_buflist_append_segment(&m->bl_srv_to_web, sbba->buf, sbba->len) < 0)
- lwsl_notice("%s: failed to store buflist segment\n", __func__);
-}
-
/*
* We will copy the buflist bl on to every sai-web client connected to our
* sai-server server, then empty bl.
@@ -120,18 +105,58 @@ _sais_websrv_broadcast_buflist(struct lws_ss_handle *h, void *arg)
void
sais_websrv_broadcast_buflist(struct lws_ss_handle *hsrv, struct lws_buflist **bl)
{
- while (*bl) {
- struct sai_bl_args sbba;
+ size_t total = 0, max_len;
+ uint8_t *flat;
+ lws_wsmsg_info_t info;
- sbba.len = lws_buflist_next_segment_len(bl, &sbba.buf);
+ if (!bl || !*bl)
+ return;
+
+ max_len = lws_buflist_total_len(bl);
+ if (!max_len) {
+ lws_buflist_destroy_all_segments(bl);
+ return;
+ }
- lws_ss_server_foreach_client(hsrv,
- _sais_websrv_broadcast_buflist,
- (void *)&sbba);
+ /*
+ * We flatten it into a single contiguous buffer so we can broadcast
+ * it as a single SOM | EOM message, which prevents other messages
+ * getting interleaved in the middle of it in the upstream buflist.
+ */
- lws_buflist_use_segment(bl, sbba.len);
+ flat = malloc(LWS_PRE + sizeof(int) + max_len);
+ if (!flat) {
+ lwsl_err("%s: OOM\n", __func__);
+ lws_buflist_destroy_all_segments(bl);
+ return;
+ }
+
+ while (*bl) {
+ uint8_t *frag;
+ size_t flen = lws_buflist_next_segment_len(bl, &frag);
+
+ if (flen > sizeof(int)) {
+ memcpy(flat + LWS_PRE + sizeof(int) + total,
+ frag + sizeof(int), flen - sizeof(int));
+ total += flen - sizeof(int);
+ }
+ lws_buflist_use_segment(bl, flen);
+ }
+
+ if (!total) {
+ free(flat);
+ return;
}
- lws_buflist_destroy_all_segments(bl);
+
+ memset(&info, 0, sizeof(info));
+ info.private_source_idx = SAI_WEBSRV_PB__PROXIED_FROM_BUILDER_LR;
+ info.buf = flat + LWS_PRE + sizeof(int);
+ info.len = total;
+ info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM;
+
+ lws_ss_server_foreach_client(hsrv, _sais_websrv_broadcast, &info);
+
+ free(flat);
}
@@ -322,6 +347,12 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid)
return SAI_DB_RESULT_ERROR;
}
+ /*
+ * Recompute startable task platforms and broadcast to all sai-power,
+ * after there has been a change in tasks
+ */
+ sais_platforms_with_tasks_pending(vhd);
+
return SAI_DB_RESULT_OK;
}
diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c
index de48d52..1bf0ee4 100644
--- a/src/server/s-ws-builder.c
+++ b/src/server/s-ws-builder.c
@@ -28,7 +28,7 @@
#include <libwebsockets.h>
#include <string.h>
-#include <signal.h>
+
#include <assert.h>
#include <time.h>
@@ -550,6 +550,10 @@ sais_process_rej(struct vhd *vhd, struct pss *pss,
break;
/* leave the uuid listed as inflight until step completed */
+ if (sais_is_task_inflight(vhd, sp, rej->task_uuid, &ul)) {
+ // lwsl_notice("%s: setting inflight started to 1 for %s\n", __func__, rej->task_uuid);
+ ul->started = 1;
+ }
break;
case SAI_TASK_REASON_DUPE:
diff --git a/src/server/s-ws-web.c b/src/server/s-ws-web.c
index bc20a07..d0752f5 100644
--- a/src/server/s-ws-web.c
+++ b/src/server/s-ws-web.c
@@ -100,7 +100,11 @@ static const lws_struct_map_t lsm_schema_json_map[] = {
LSM_SCHEMA (sai_pcon_control_t, NULL, lsm_pcon_control,
/* shares struct */ "com.warmcat.sai.pcon_control"),
LSM_SCHEMA (sai_browse_rx_taskinfo_t, NULL, lsm_browser_taskinfo,
- "com.warmcat.sai.taskinfo")
+ "com.warmcat.sai.taskinfo"),
+ LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset,
+ /* shares struct */ "com.warmcat.sai.taskpause"),
+ LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset,
+ /* shares struct */ "com.warmcat.sai.taskresume")
};
enum {
@@ -115,6 +119,8 @@ enum {
SAIS_WS_WEBSRV_RX_STAY,
SAIS_WS_WEBSRV_RX_PCON_CONTROL,
SAIS_WS_WEBSRV_RX_TASKINFO,
+ SAIS_WS_WEBSRV_RX_TASKPAUSE,
+ SAIS_WS_WEBSRV_RX_TASKRESUME,
};
static int
@@ -442,6 +448,28 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
lwsl_ss_err(m->ss, "taskreset failed");
break;
+ case SAIS_WS_WEBSRV_RX_TASKPAUSE:
+ ei = (sai_browse_rx_evinfo_t *)a.dest;
+ if (sais_validate_id(ei->event_hash, SAI_TASKID_LEN))
+ goto soft_error;
+
+ lwsl_ss_warn(m->ss, "SAIS_WS_WEBSRV_RX_TASKPAUSE: %s: received", ei->event_hash);
+ if (sais_task_pause(m->vhd, ei->event_hash))
+ lwsl_ss_err(m->ss, "taskpause failed");
+ break;
+
+ case SAIS_WS_WEBSRV_RX_TASKRESUME:
+ ei = (sai_browse_rx_evinfo_t *)a.dest;
+ if (sais_validate_id(ei->event_hash, SAI_TASKID_LEN))
+ goto soft_error;
+
+ lwsl_ss_warn(m->ss, "SAIS_WS_WEBSRV_RX_TASKRESUME: %s: received", ei->event_hash);
+ if (sais_set_task_state(m->vhd, ei->event_hash, SAIES_WAITING, 0, 0))
+ lwsl_ss_err(m->ss, "taskresume failed");
+ else
+ sais_platforms_with_tasks_pending(m->vhd);
+ break;
+
case SAIS_WS_WEBSRV_RX_TASKREBUILDLASTSTEP:
ei = (sai_browse_rx_evinfo_t *)a.dest;
if (sais_validate_id(ei->event_hash, SAI_TASKID_LEN))
diff --git a/src/web/w-comms.c b/src/web/w-comms.c
index 13d5efb..bc0931f 100644
--- a/src/web/w-comms.c
+++ b/src/web/w-comms.c
@@ -918,21 +918,40 @@ clean_spa:
break;
{
- int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&pss->raw_tx), depi = *pi;
char som, eom, rb[1200];
int used, final = 1;
size_t fsl = lws_buflist_next_segment_len(&pss->raw_tx, NULL);
- /* this is the only buflist user on pss->raw_tx */
- used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, sizeof(rb), &som, &eom);
- if (!used)
+ /*
+ * Each segment has a header containing the flags.
+ * We MUST only read it if we are at the start of the segment.
+ * If we are mid-segment, we use the cached flags.
+ */
+ if (lws_buflist_get_frag_start_or_NULL(&pss->raw_tx)) {
+ /* This is just a peek to see if we HAVE a segment */
+ int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&pss->raw_tx);
+ int flags = *pi;
+
+ /*
+ * fragment_use sets 'som' to true if we are at
+ * the segment start.
+ */
+ used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, sizeof(rb), &som, &eom);
+ if (!used)
+ return 0;
+
+ if (som)
+ pss->segment_flags = flags;
+ } else
return 0;
- if (used < (int)fsl || (depi & LWS_WRITE_NO_FIN))
+
+ if (used < (int)fsl || (pss->segment_flags & LWS_WRITE_NO_FIN))
final = 0;
if (lws_write(pss->wsi, (uint8_t *)rb + ((size_t)som * sizeof(int)),
(size_t)used - ((size_t)som * sizeof(int)),
- (lws_ws_sending_multifragment(pss->wsi) ? LWS_WRITE_CONTINUATION : LWS_WRITE_TEXT) |
+ (lws_ws_sending_multifragment(pss->wsi) ?
+ LWS_WRITE_CONTINUATION : LWS_WRITE_TEXT) |
(!final * LWS_WRITE_NO_FIN)) < 0) {
lwsl_wsi_err(pss->wsi, "attempt to write %d failed", (int)used - (int)sizeof(int));
diff --git a/src/web/w-private.h b/src/web/w-private.h
index 7717cd8..25a47f9 100644
--- a/src/web/w-private.h
+++ b/src/web/w-private.h
@@ -98,6 +98,7 @@ struct pss {
int log_cache_size;
int authorized;
int specificity;
+ int segment_flags;
unsigned int js_api_version;
unsigned long expiry_unix_time;
@@ -134,6 +135,9 @@ struct vhd {
struct lws_dll2_owner builders_owner;
struct lwsac *builders;
+ struct lws_dll2_owner pcons_owner;
+ struct lwsac *pcons;
+
/* our keys */
struct lws_jwk jwt_jwk_auth;
char jwt_auth_alg[16];
@@ -239,5 +243,7 @@ saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss);
int
saiw_browser_broadcast_queue_builders(struct vhd *vhd, struct pss *pss);
+int
+saiw_browser_broadcast_queue_pcons(struct vhd *vhd, struct pss *pss);
diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c
index bc3f33a..9ad2c23 100644
--- a/src/web/w-ws-browser.c
+++ b/src/web/w-ws-browser.c
@@ -821,6 +821,12 @@ saiw_broadcast_logs_batch(struct vhd *vhd, struct pss *pss)
if (!pss->subs_list.owner)
return 0;
+ if (lws_buflist_total_len(&pss->raw_tx) > 100 * 1024) {
+ lws_sul_schedule(vhd->context, 0, &pss->sul_logcache,
+ saiw_retry_logs, 250 * LWS_US_PER_MS);
+ return 0;
+ }
+
/*
* For efficiency, let's try to grab the next 100 at
* once from sqlite and work our way through sending
@@ -1027,41 +1033,54 @@ saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss)
lwsl_err("%s: json ser fail\n", __func__);
return 1;
}
+ if (lws_ptr_diff_size_t(end, p) < 128) {
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
+ p = start;
+ }
+
if (subsequent)
*p++ = ',';
subsequent = 1;
p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "{\"e\":");
- if (lws_ptr_diff_size_t(end, p) < 256) {
- saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
- lws_ptr_diff_size_t(p, start),
- lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
- p = start;
- }
- n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w);
- lws_struct_json_serialize_destroy(&js);
- switch (n) {
- case LSJS_RESULT_ERROR:
- lwsl_err("%s: json ser error\n", __func__);
- return 1;
+ do {
+ n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w);
+ switch (n) {
+ case LSJS_RESULT_ERROR:
+ lwsl_err("%s: json ser error\n", __func__);
+ lws_struct_json_serialize_destroy(&js);
+ return 1;
- case LSJS_RESULT_FINISH:
- case LSJS_RESULT_CONTINUE:
- p += w;
- task_index = 0;
- p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), ", \"t\":[");
- break;
- }
+ case LSJS_RESULT_FINISH:
+ lws_struct_json_serialize_destroy(&js);
+ p += w;
+ break;
- if (lws_ptr_diff_size_t(end, p) < 2560) {
+ case LSJS_RESULT_CONTINUE:
+ p += w;
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
+ p = start;
+ break;
+ }
+ } while (n == LSJS_RESULT_CONTINUE);
+
+ if (lws_ptr_diff_size_t(end, p) < 128) {
saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
lws_ptr_diff_size_t(p, start),
lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
p = start;
}
+ task_index = 0;
+ p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), ", \"t\":[");
+
+
/*
* Enumerate the tasks associated with this event...
*/
@@ -1124,23 +1143,39 @@ saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss)
LWS_ARRAY_SIZE(lsm_schema_json_map_task), 0, t);
t->build[0] = '\0';
- n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w);
- lws_struct_json_serialize_destroy(&js);
- lwsac_free(&task_ac);
- p += w;
- if (lws_ptr_diff_size_t(end, p) < 2560) {
- saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
- lws_ptr_diff_size_t(p, start),
- lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
- p = start;
- }
+ do {
+ n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w);
+ switch (n) {
+ case LSJS_RESULT_FINISH:
+ lws_struct_json_serialize_destroy(&js);
+ p += w;
+ break;
+
+ case LSJS_RESULT_CONTINUE:
+ p += w;
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
+ p = start;
+ break;
+ }
+ } while (n == LSJS_RESULT_CONTINUE);
+
+ lwsac_free(&task_ac);
task_index++;
} while (1);
/* none left to do, go back up a level */
+ if (lws_ptr_diff_size_t(end, p) < 128) {
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
+ p = start;
+ }
+
p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}");
if (pss->specificity)
@@ -1153,6 +1188,13 @@ saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss)
}
so_finish:
+ if (lws_ptr_diff_size_t(end, p) < 16) {
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0));
+ p = start;
+ }
+
p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}");
saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
@@ -1163,8 +1205,58 @@ so_finish:
}
int
+saiw_browser_broadcast_queue_pcons(struct vhd *vhd, struct pss *pss)
+{
+ char buf[4096 + LWS_PRE], *start = buf + LWS_PRE, *p = start,
+ *end = buf + sizeof(buf);
+ lws_struct_serialize_t *js;
+ sai_power_managed_builders_t pmb;
+ lws_struct_json_serialize_result_t r;
+ size_t w;
+ char fi = 1;
+
+ if (!vhd || !vhd->pcons)
+ return 0;
+
+ memset(&pmb, 0, sizeof(pmb));
+ pmb.power_controllers = vhd->pcons_owner;
+
+ js = lws_struct_json_serialize_create(
+ lsm_schema_power_managed_builders,
+ LWS_ARRAY_SIZE(lsm_schema_power_managed_builders),
+ 0, &pmb);
+ if (!js)
+ return 1;
+
+ do {
+ r = lws_struct_json_serialize(js, (uint8_t *)p,
+ lws_ptr_diff_size_t(end, p), &w);
+ p += w;
+
+ switch (r) {
+ case LSJS_RESULT_FINISH:
+ case LSJS_RESULT_CONTINUE:
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, fi, r == LSJS_RESULT_FINISH));
+ fi = 0;
+ p = start;
+ break;
+ case LSJS_RESULT_ERROR:
+ lws_struct_json_serialize_destroy(&js);
+ return 1;
+ }
+ } while (r == LSJS_RESULT_CONTINUE);
+
+ lws_struct_json_serialize_destroy(&js);
+
+ return 0;
+}
+
+int
saiw_browser_broadcast_queue_builders(struct vhd *vhd, struct pss *pss)
{
+ saiw_browser_broadcast_queue_pcons(vhd, pss);
char buf[4096 + LWS_PRE], *start = buf + LWS_PRE, *p = start,
*end = buf + sizeof(buf);
lws_struct_serialize_t *js;
@@ -1192,6 +1284,10 @@ saiw_browser_broadcast_queue_builders(struct vhd *vhd, struct pss *pss)
while (walk) {
sai_plat_t *b = lws_container_of(walk, sai_plat_t, sai_plat_list);
+ lws_struct_json_serialize_result_t r;
+ char start_of_this_builder = 1;
+
+ lwsl_notice("%s: processing builder '%s' (online %d)\n", __func__, b->name, b->online);
js = lws_struct_json_serialize_create(
lsm_schema_map_plat_simple,
@@ -1201,38 +1297,63 @@ saiw_browser_broadcast_queue_builders(struct vhd *vhd, struct pss *pss)
lwsac_unreference(&vhd->builders);
return 1;
}
- if (subsequent)
- *p++ = ',';
- subsequent = 1;
- switch (lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w)) {
- case LSJS_RESULT_ERROR:
- lws_struct_json_serialize_destroy(&js);
- return 1;
+ do {
+ if (subsequent && start_of_this_builder) {
+ *p++ = ',';
+ start_of_this_builder = 0;
+ }
- case LSJS_RESULT_FINISH:
- lws_struct_json_serialize_destroy(&js);
- /* fallthru */
- case LSJS_RESULT_CONTINUE:
+ r = lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p) - 2, &w);
p += w;
- walk = walk->next;
- break;
- }
- if (lws_ptr_diff_size_t(end, p) < 256) {
- saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ switch (r) {
+ case LSJS_RESULT_CONTINUE:
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
lws_ptr_diff_size_t(p, start),
lws_write_ws_flags(LWS_WRITE_TEXT, fi, 0));
+ fi = 0;
+ p = start;
+ break;
+ case LSJS_RESULT_ERROR:
+ lws_struct_json_serialize_destroy(&js);
+ lwsac_unreference(&vhd->builders);
+ return 1;
+ case LSJS_RESULT_FINISH:
+ lws_struct_json_serialize_destroy(&js);
+ break;
+ }
+ } while (r == LSJS_RESULT_CONTINUE);
+
+ subsequent = 1;
+ walk = walk->next;
+
+ if (walk && lws_ptr_diff_size_t(end, p) < 512) {
+ /* No room for another builder, fragment now */
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, fi, 0));
fi = 0;
p = start;
}
}
- p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}");
+ if (lws_ptr_diff_size_t(end, p) < 16) {
+ saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, fi, 0));
+ fi = 0;
+ p = start;
+ }
+
+ p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), " \n]}");
saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start,
- lws_ptr_diff_size_t(p, start),
- lws_write_ws_flags(LWS_WRITE_TEXT, fi, 1));
+ lws_ptr_diff_size_t(p, start),
+ lws_write_ws_flags(LWS_WRITE_TEXT, fi, 1));
+
+ lwsac_unreference(&vhd->builders);
+
return 0;
}
diff --git a/src/web/w-ws-server.c b/src/web/w-ws-server.c
index a1d3e82..2261a9e 100644
--- a/src/web/w-ws-server.c
+++ b/src/web/w-ws-server.c
@@ -87,13 +87,14 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
{
saiw_websrv_t *m = (saiw_websrv_t *)userobj;
struct vhd *vhd = (struct vhd *)m->opaque_data;
- sai_browse_rx_evinfo_t *ei;
- int n;
+ int n, is_start = (flags & LWSSS_FLAG_SOM);
+ const uint8_t *p = buf;
+ size_t rem = len;
// lwsl_ss_warn(m->ss, "%s: len %d, flags %d\n", __func__, (int)len, flags);
// lwsl_hexdump_notice(buf, len);
- if (flags & LWSSS_FLAG_SOM) {
+ if (is_start) {
/* First frag of a new message. Clear old parse results and init */
lwsac_free(&m->a.ac);
memset(&m->a, 0, sizeof(m->a));
@@ -106,189 +107,178 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
lws_struct_json_init_parse(&m->ctx, NULL, &m->a);
}
- // fprintf(stderr, "%s: rx: %.*s\n", __func__, (int)len, buf);
+ while (rem > 0) {
+ n = lejp_parse(&m->ctx, (uint8_t *)p, (int)rem);
- n = lejp_parse(&m->ctx, (uint8_t *)buf, (int)len);
-
- /* Check for fatal error OR completion without an object */
- if (n < 0 && n != LEJP_CONTINUE) {
- lwsl_notice("%s: srv->web JSON decode failed '%s' (ssflags %d)\n",
- __func__, lejp_error_to_string(n), flags);
- lwsl_hexdump_notice(buf, len);
- goto cleanup_and_disconnect;
- }
+ /* Check for fatal error OR completion without an object */
+ if (n < 0 && n != LEJP_CONTINUE) {
+ lwsl_notice("%s: srv->web JSON decode failed '%s' (ssflags %d)\n",
+ __func__, lejp_error_to_string(n), flags);
+ lwsl_hexdump_notice(p, rem);
+ goto cleanup_and_disconnect;
+ }
- 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:
- case SAIS_WS_WEBSRV_RX_TASKACTIVITY:
- case SAIS_WS_WEBSRV_RX_SAI_BUILDERS:
- case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS:
- case SAIS_WS_WEBSRV_RX_PCON_ENERGY:
- saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len,
- lws_write_ws_flags(LWS_WRITE_TEXT,
- flags & LWSSS_FLAG_SOM,
- flags & LWSSS_FLAG_EOM));
- break;
- default:
- lwsl_err("%s: SWALLOWING %.*s\n", __func__, (int)len, buf);
- break;
+ 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:
+ case SAIS_WS_WEBSRV_RX_TASKACTIVITY:
+ case SAIS_WS_WEBSRV_RX_SAI_BUILDERS:
+ case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS:
+ case SAIS_WS_WEBSRV_RX_PCON_ENERGY:
+ saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, p, rem,
+ lws_write_ws_flags(LWS_WRITE_TEXT,
+ is_start,
+ 0)); /* Not EOM */
+ break;
+ default:
+ // lwsl_err("%s: SWALLOWING %.*s\n", __func__, (int)len, buf);
+ break;
+ }
+
+ return 0;
}
- return 0;
- } else {
+ /* We have a completed message */
+ size_t consumed = rem - (size_t)n;
+
+ sai_browse_rx_evinfo_t *ei;
+
switch (m->a.top_schema_index) {
case SAIS_WS_WEBSRV_RX_TASKCHANGE:
case SAIS_WS_WEBSRV_RX_EVENTCHANGE:
case SAIS_WS_WEBSRV_RX_SAI_BUILDERS:
case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS:
case SAIS_WS_WEBSRV_RX_PCON_ENERGY:
- saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len,
+ case SAIS_WS_WEBSRV_RX_LOADREPORT:
+ case SAIS_WS_WEBSRV_RX_TASKACTIVITY:
+ saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, p, consumed,
lws_write_ws_flags(LWS_WRITE_TEXT,
- flags & LWSSS_FLAG_SOM,
- flags & LWSSS_FLAG_EOM));
+ is_start,
+ 1)); /* Force EOM */
break;
}
- // lwsl_err("%s: proxying %.*s\n", __func__, (int)len, buf);
- }
-
- /*
- * 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;
- lwsl_notice("%s: TASKCHANGE %s\n", __func__, ei->event_hash);
- saiw_browsers_task_state_change(vhd, ei->event_hash);
- break;
-
- case SAIS_WS_WEBSRV_RX_EVENTCHANGE:
- ei = (sai_browse_rx_evinfo_t *)m->a.dest;
- lwsl_notice("%s: EVENTCHANGE %s\n", __func__, ei->event_hash);
- saiw_event_state_change(vhd, ei->event_hash);
- break;
+ /*
+ * 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;
+ }
- 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 */
+ switch (m->a.top_schema_index) {
- /* 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 *sp = lws_container_of(p, sai_plat_t, sai_plat_list);
+ case SAIS_WS_WEBSRV_RX_TASKCHANGE:
+ ei = (sai_browse_rx_evinfo_t *)m->a.dest;
+ lwsl_notice("%s: TASKCHANGE %s\n", __func__, ei->event_hash);
+ saiw_browsers_task_state_change(vhd, ei->event_hash);
+ break;
- lws_dll2_remove(&sp->sai_plat_list);
- lws_dll2_add_tail(&sp->sai_plat_list, &vhd->builders_owner);
- } lws_end_foreach_dll_safe(p, p1);
+ case SAIS_WS_WEBSRV_RX_EVENTCHANGE:
+ ei = (sai_browse_rx_evinfo_t *)m->a.dest;
+ lwsl_notice("%s: EVENTCHANGE %s\n", __func__, ei->event_hash);
+ saiw_event_state_change(vhd, ei->event_hash);
+ break;
- /* 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);
+ 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 */
+
+ /* 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 *sp = lws_container_of(p, sai_plat_t, sai_plat_list);
+
+ lws_dll2_remove(&sp->sai_plat_list);
+ lws_dll2_add_tail(&sp->sai_plat_list, &vhd->builders_owner);
+ } lws_end_foreach_dll_safe(p, p1);
+
+ /* 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);
+
+ saiw_browser_broadcast_queue_builders(pss->vhd, pss);
+ } lws_end_foreach_dll(p);
+ break;
- saiw_browser_broadcast_queue_builders(pss->vhd, pss);
- } lws_end_foreach_dll(p);
- break;
+ case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS:
+ lwsac_free(&vhd->pcons);
+ lws_dll2_owner_clear(&vhd->pcons_owner);
+ vhd->pcons = m->a.ac;
+ m->a.ac = NULL; /* The vhd now owns this memory */
+
+ /* Move the parsed objects to the vhd's list */
+ lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1,
+ ((sai_power_managed_builders_t *)m->a.dest)->power_controllers.head) {
+ sai_power_controller_t *pc = lws_container_of(p, sai_power_controller_t, list);
+
+ lws_dll2_remove(&pc->list);
+ lws_dll2_add_tail(&pc->list, &vhd->pcons_owner);
+ } lws_end_foreach_dll_safe(p, p1);
+
+ /* 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);
+
+ saiw_browser_broadcast_queue_builders(pss->vhd, pss);
+ } 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);
+ 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_browser_queue_overview(pss->vhd, pss);
- } lws_end_foreach_dll(p);
- break;
+ saiw_browser_queue_overview(pss->vhd, pss);
+ } lws_end_foreach_dll(p);
+ break;
- case SAIS_WS_WEBSRV_RX_TASKLOGS:
- ei = (sai_browse_rx_evinfo_t *)m->a.dest;
- 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))
- saiw_broadcast_logs_batch(vhd, pss);
- } lws_end_foreach_dll(p);
- break;
+ case SAIS_WS_WEBSRV_RX_TASKLOGS:
+ ei = (sai_browse_rx_evinfo_t *)m->a.dest;
+ 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))
+ saiw_broadcast_logs_batch(vhd, pss);
+ } lws_end_foreach_dll(p);
+ break;
+ }
- case SAIS_WS_WEBSRV_RX_LOADREPORT:
- // lwsl_notice("%s: ^^^^^^^^^^^^^^ SAIS_WS_WEBSRV_RX_LOADREPORT forwarding to browser\n", __func__);
- // lwsl_hexdump_notice(buf, len);
- saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n,
- lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM));
- break;
- case SAIS_WS_WEBSRV_RX_TASKACTIVITY:
- saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n,
- lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM));
- break;
- case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS:
- /* Just forward to browser */
- /* We already forwarded fragments in the LEJP_CONTINUE block?
- * Wait, in the LEJP_CONTINUE block above, we handle partial forwarding.
- * But if it arrives in one chunk (n >= 0 immediately), we must forward it here.
- * The existing logic for other types does this.
- * Wait, the LOADREPORT/TASKACTIVITY/BUILDERS cases below just call broadcast.
- * But they use 'len - n' which assumes n is the length processed?
- * Lejp returns LEJP_CONTINUE or a positive integer representing the number of bytes used?
- * No, lejp_parse returns LEJP_CONTINUE or 0 (success) or error code.
- *
- * Wait, lejp_parse returns a negative error code on failure.
- * On success, it returns 0.
- *
- * The code `len - (unsigned int)n` suggests n is bytes consumed?
- * Ah, older versions of lejp/lws might behave differently.
- * But if n == 0 (success), `len - 0` = len.
- * So it broadcasts the whole buffer.
- *
- * The `else` block of `if (n == LEJP_CONTINUE)` handles the `flags & EOM` case for fragments?
- * No, `if (n == LEJP_CONTINUE)` handles forwarding intermediate fragments.
- * The `else` block handles the FINAL fragment (where n == 0).
- * But wait, the `else` block logic is:
- *
- * ```c
- * } else {
- * switch (m->a.top_schema_index) {
- * case SAIS_WS_WEBSRV_RX_TASKCHANGE: ... broadcast ...
- * }
- * }
- * ```
- *
- * It seems correct. I added `SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS` to the `else` block switch as well.
- * But I should check if I need to add a case in the final switch (where parsing is complete).
- *
- * In the final switch (m->a.top_schema_index), `SAIS_WS_WEBSRV_RX_LOADREPORT` etc are handled.
- * I should add `SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS` there too if I want to ensure the final chunk is sent if it wasn't handled by the `else` block?
- *
- * Actually, look at `SAIS_WS_WEBSRV_RX_SAI_BUILDERS` in the final switch. It does complex logic (updating vhd state).
- * `com.warmcat.sai.power_managed_builders` is just passthrough to the browser, sai-web doesn't need to statefully track PCONs (browsers do).
- * So treating it like LOADREPORT (passthrough) is correct.
+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.
*/
- saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n,
- lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM));
- break;
- case SAIS_WS_WEBSRV_RX_PCON_ENERGY:
- saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n,
- lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM));
- break;
+ lwsac_free(&m->a.ac);
+
+ /* Advance to next part of buffer */
+ p += consumed;
+ rem = (size_t)n; // unused bytes
+
+ if (rem > 0) {
+ /* Prepare for next message */
+ 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_st[1] = lsm_schema_json_map;
+ m->a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_json_map);
+ m->a.ac_block_size = 4096;
+
+ lws_struct_json_init_parse(&m->ctx, NULL, &m->a);
+
+ /* Subsequent messages in same packet are always new Starts */
+ is_start = 1;
+ }
}
-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: