Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / src / web / w-private.h
Author[]Andy Green <andy@warmcat.com> 2026-03-19 11:09 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-03-26 09:38 UTC
Treeeab625c77addd7d339aedf37e10e971f6fca01cc   Raw Patch
 
logs: fix update in browser
logs: fix update in browser
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 += "&#9633;"; 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:
Page fetched 0s ago, creation time: 19ms (vhost etag hits: 0%, cache hits: 0%)