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 / assets / tc-noptmsvc.svg
Author[]Andy Green <andy@warmcat.com> 2025-11-05 13:08 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-11-06 08:24 UTC
Tree6222ff98e30fdb7f099a6bfe7c970f20e9c203e3   Raw Patch
 
plat-rebuild
plat-rebuild

Various rewrites / refactors / cleaning / fixes to get us out of the
problems introduced mainly by AI commits.

This gets us mainly back to normal operation using the full set of
builders.
diff --git a/assets/sai.css b/assets/sai.css index d1d2470..c3b8e64 100644 --- a/assets/sai.css +++ b/assets/sai.css @@ -603,6 +603,9 @@ td.tn { font-weight: normal; font-size: 7pt; padding-left:0px; + white-space: nowrap; + overflow: hidden; + text-overflow: ellipsis; } table.nomar { diff --git a/assets/sai.js b/assets/sai.js index 74ba59d..c0c6e72 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -536,7 +536,7 @@ function renderSpreadsheet(tasks) { html += '<tr>' + `<td>` + s1 + `</td>` + `<td>${agify(now_ut, task.started)} ago</td>` + - `<td><a href="?task=${hsanitize(task.task_uuid)}">${hsanitize(task.task_name)}</a></td>` + + `<td><a href="index.html?task=${hsanitize(task.task_uuid)}">${hsanitize(task.task_name)}</a></td>` + '</tr>'; } diff --git a/src/builder/CMakeLists.txt b/src/builder/CMakeLists.txt index f98714b..1b1f9a3 100644 --- a/src/builder/CMakeLists.txt +++ b/src/builder/CMakeLists.txt @@ -4,7 +4,7 @@ set(CPACK_DEBIAN_BUILDER_PACKAGE_NAME "sai-builder") set(SRCS b-sai.c b-conf.c - b-comms.c + b-ws-server.c b-nspawn.c b-task.c b-artifacts.c @@ -16,6 +16,7 @@ set(SRCS b-deletion.c b-power.c ../common/c-utils.c + ../common/struct-metadata.c ) set(requirements 1) diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c deleted file mode 100644 index cbd0442..0000000 --- a/src/builder/b-comms.c +++ /dev/null @@ -1,663 +0,0 @@ -/* - * sai-builder com-warmcat-sai client protocol implementation - * - * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - */ - -#include <libwebsockets.h> -#include <string.h> -#include <signal.h> - -#include "b-private.h" - -extern struct lws_spawn_piped *lsp_suspender; - -#include "../common/struct-metadata.c" - -static const lws_struct_map_t lsm_schema_json_loadreport[] = { - LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, "com.warmcat.sai.loadreport"), -}; - -static const lws_struct_map_t lsm_viewerstate_members[] = { - LSM_UNSIGNED(sai_viewer_state_t, viewers, "viewers"), -}; - -const lws_struct_map_t lsm_schema_map_m_to_b[] = { - LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), - LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, "com.warmcat.sai.taskcan"), - LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewerstate_members, - "com.warmcat.sai.viewerstate"), - LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, "com-warmcat-sai-resource"), - LSM_SCHEMA (sai_rebuild_t, NULL, lsm_rebuild, "com.warmcat.sai.rebuild") -}; - -enum { - SAIB_RX_TASK_ALLOCATION, - SAIB_RX_TASK_CANCEL, - SAIB_RX_VIEWERSTATE, - SAIB_RX_RESOURCE_REPLY, - SAIB_RX_REBUILD -}; - -/* - * This is the only path to send things from builder->server. - * - * It will copy the incoming buffer fragment into a buflist in order. So you - * should dump all your fragments for a message in here one after the other - * and the message will go out uninterrupted. Having this as the only tx path - * allows us to guarantee we won't interrupt the fragment sequencing. - * - * The fragment sizing does not have to be related to ss usage sizing, it can - * be larger and it will be used from the buflist according to what SS wants. - */ - -int -saib_srv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)lws_ss_to_user_object(h); - unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); - - *pi = ss_flags; - - // lwsl_ss_notice(h, "Queuing builder -> sai-server"); - // lwsl_hexdump_notice(buf, len); - - if (lws_buflist_append_segment(&spm->bl_to_srv, (uint8_t *)buf - sizeof(int), - len + sizeof(int)) < 0) - lwsl_ss_err(h, "failed to append"); /* still ask to drain */ - - if (lws_ss_request_tx(h)) - lwsl_ss_err(h, "failed to request tx"); - - return 0; -} - -int -saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, - const lws_struct_map_t *map, - size_t map_entries, void *object) -{ - lws_struct_serialize_t *js; - unsigned int ssf = LWSSS_FLAG_SOM; - uint8_t buf[1024 + LWS_PRE]; - size_t w = 0; - - js = lws_struct_json_serialize_create(map, map_entries, 0, object); - if (!js) { - lwsl_warn("%s: failed to serialize\n", __func__); - return -1; - } - - do { - switch (lws_struct_json_serialize(js, buf + LWS_PRE, - sizeof(buf) - LWS_PRE, &w)) { - case LSJS_RESULT_CONTINUE: - break; - case LSJS_RESULT_FINISH: - ssf |= LWSSS_FLAG_EOM; - break; - case LSJS_RESULT_ERROR: - lwsl_warn("%s: serialization failed\n", __func__); - return -1; - } - - if (saib_srv_queue_tx(h, buf + LWS_PRE, w, ssf)) - return -1; - - ssf &= ~((unsigned int)LWSSS_FLAG_SOM); - } while (!(ssf & LWSSS_FLAG_EOM)); - - lws_struct_json_serialize_destroy(&js); - - return 0; -} - -static lws_ss_state_return_t -saib_m_rx(void *userobj, const uint8_t *in, size_t len, int flags) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)userobj; - sai_plat_t *sp = NULL; - sai_resource_t *reso; - struct lejp_ctx ctx; - lws_struct_args_t a; - sai_cancel_t *can; - sai_rebuild_t *reb; - int m; - - /* - * use the schema name on the incoming JSON to decide what kind of - * structure to instantiate - */ - - memset(&a, 0, sizeof(a)); - a.map_st[0] = lsm_schema_map_m_to_b; - a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_m_to_b); - a.ac_block_size = 512; - -// lwsl_hexdump_warn(in, len); - - lws_struct_json_init_parse(&ctx, NULL, &a); - m = lejp_parse(&ctx, (uint8_t *)in, (int)len); - if (m < 0) { - lwsl_hexdump_err(in, len); - lwsl_err("%s: builder rx JSON decode failed '%s'\n", - __func__, lejp_error_to_string(m)); - return m; - } - - if (!a.dest) { - lwsac_free(&a.ac); - return LWSSSSRET_OK; - } - - switch (a.top_schema_index) { - - case SAIB_RX_TASK_ALLOCATION: - if (saib_consider_allocating_task(spm, &a, in, len, flags)) - break; - - break; - - case SAIB_RX_TASK_CANCEL: - - can = (sai_cancel_t *)a.dest; - - lwsl_notice("%s: received task cancel for %s\n", __func__, can->task_uuid); - - lws_start_foreach_dll_safe(struct lws_dll2 *, mp, mp1, - builder.sai_plat_owner.head) { - struct sai_plat *sp = lws_container_of(mp, struct sai_plat, - sai_plat_list); - - lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, - sp->nspawn_owner.head) { - struct sai_nspawn *ns = lws_container_of(p, - struct sai_nspawn, list); - - if (ns->task && - !strcmp(can->task_uuid, ns->task->uuid)) { - lwsl_notice("%s: trying to cancel %s\n", - __func__, can->task_uuid); - - /* - * We're going to send a few signals - * at 500ms intervals - */ - ns->user_cancel = 1; - ns->term_budget = 5; - - lws_sul_schedule(ns->builder->context, 0, - &ns->sul_task_cancel, - saib_sul_task_cancel, 1); - } - - } lws_end_foreach_dll_safe(p, p1); - - } lws_end_foreach_dll_safe(mp, mp1); - break; - - case SAIB_RX_VIEWERSTATE: - { - sai_viewer_state_t *vs = (sai_viewer_state_t *)a.dest; - char any_busy = 0; - - lwsl_notice("Received viewer state update: %u viewers\n", vs->viewers); - - spm->viewer_count = vs->viewers; - - if (!vs->viewers) { - lwsl_notice("%s: VIEWERSTATE: no viewers -> no load reports\n", __func__); - lws_sul_cancel(&spm->sul_load_report); - break; - } - - /* are there any busy instances */ - - lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, - builder.sai_plat_owner.head) { - sp = lws_container_of(d, sai_plat_t, sai_plat_list); - - lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { - struct sai_nspawn *ns = lws_container_of(d, struct sai_nspawn, list); - - if (ns->state == NSSTATE_EXECUTING_STEPS) - any_busy = 1; - - } lws_end_foreach_dll(d); - } lws_end_foreach_dll_safe(d, d1); - - if (!any_busy) { - lwsl_notice("%s: VIEWERSTATE: no busy instances -> no load reports\n", __func__); - - lws_sul_cancel(&spm->sul_load_report); - break; - } - - /* At least one viewer, start reporting */ - lwsl_notice("%s: VIEWERSTATE: viewers + busy instances -> load reports\n", __func__); - - lws_sul_schedule(builder.context, 0, &spm->sul_load_report, - saib_sul_load_report_cb, 1); - } - break; - - case SAIB_RX_RESOURCE_REPLY: - reso = (sai_resource_t *)a.dest; - - lwsl_notice("%s: RESOURCE_REPLY: cookie %s\n", - __func__, reso->cookie); - - saib_handle_resource_result(spm, (const char *)in, len); - break; - - case SAIB_RX_REBUILD: - reb = (sai_rebuild_t *)a.dest; - - lwsl_notice("%s: REBUILD: %s\n", __func__, reb->builder_name); - - if (suspender_exists) { - uint8_t b = 3; - int fd = saib_suspender_get_pipe(); - - if (write(fd, &b, 1) != 1) - lwsl_err("%s: Failed to write to suspender\n", - __func__); - } - break; - - default: - break; - } - - return LWSSSSRET_OK; -} - -/* - * We cover requested tx for any instance of a platform that can takes tasks - * from the same server... it means just by coming here, no particular - * platform / sai_plat is implied... - */ - -static lws_ss_state_return_t -saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, - int *flags) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)userobj; - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&spm->bl_to_srv), depi; - char som, som1, eom, final = 1; - size_t fsl, used; - - if (!spm->bl_to_srv) { - lwsl_notice("%s: nothing to send from builder -> srv\n", __func__); - return LWSSSSRET_TX_DONT_SEND; - } - - depi = *pi; - *pi = (*pi) & (~(LWSSS_FLAG_SOM)); /* no SOM twice even on partial */ - - /* - * We can only issue *len at a time. - * - * Notice we are getting the stored flags from the START of the fragment each time. - * that means we can still see the right flags stored with the fragment, even if we - * have partially used the buflist frag and are partway through it. - * - * Ergo, only something to skip if we are at som=1. And also notice that although - * *pi will be right, after the lws_buflist..._use() api, what it points to has been - * destroyed. So we also dereference *pi into depi for use below. - */ - - fsl = lws_buflist_next_segment_len(&spm->bl_to_srv, NULL); - - lws_buflist_fragment_use(&spm->bl_to_srv, NULL, 0, &som, &eom); - if (som) { - fsl -= sizeof(int); - lws_buflist_fragment_use(&spm->bl_to_srv, buf, sizeof(int), &som1, &eom); - } - if (!(depi & LWSSS_FLAG_SOM)) - som = 0; - - used = (size_t)lws_buflist_fragment_use(&spm->bl_to_srv, (uint8_t *)buf, *len, &som1, &eom); - if (!used) - return LWSSSSRET_TX_DONT_SEND; - - if (used < fsl || !(depi & LWSSS_FLAG_EOM)) - final = 0; - - *len = used; - *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - -// lwsl_ss_notice(spm->ss, "Sending %d builder->srv: ssflags %d", (int)*len, (int)*flags); -// lwsl_hexdump_notice(buf, *len); - - if (spm->bl_to_srv) - return lws_ss_request_tx(spm->ss); - - return 0; -} - -static int -cleanup_on_ss_destroy(struct lws_dll2 *d, void *user) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)user; - sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list); - - lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, - sp->nspawn_owner.head) { - struct sai_nspawn *ns = - lws_container_of(d, struct sai_nspawn, list); - - if (ns->spm == spm) { - lwsl_warn("%s: ns->spm %p, spm %p\n", __func__, ns->spm, spm); - /* - * This pss is about to go away, make sure the ns - * can't reference it any more no matter what happens - */ - ns->spm = NULL; - } - } lws_end_foreach_dll_safe(d, d1); - - return 0; -} - -static int -cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)user; - sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list); - - lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, - sp->nspawn_owner.head) { - struct sai_nspawn *ns = lws_container_of(d, - struct sai_nspawn, list); - - if (ns->spm == spm) { - - /* - * This pss is about to go away, make sure the ns - * can't reference it any more no matter what happens - */ - - ns->spm = NULL; - - if (ns->op && ns->op->lsp) - lws_spawn_piped_kill_child_process(ns->op->lsp); - } - } lws_end_foreach_dll_safe(d, d1); - - return 0; -} - -void -saib_sul_load_report_cb(struct lws_sorted_usec_list *sul) -{ - struct sai_plat_server *spm = lws_container_of(sul, - struct sai_plat_server, sul_load_report); - char any_platform_on_this_spm_active = 0; - int n; - - /* - * This builder process may have multiple platforms, each with - * multiple instances. We report on each platform separately so the - * UI can distinguish them. - * - * This SUL is per-server-connection. We iterate all platforms and - * for each, see if it's supposed to connect to this server. - * - * To avoid spamming idle reports, we only report if the platform - * is active for this server, OR if it was active the last time we - * checked (ie, it has just become idle, so we need to send one last - * report with no active tasks to clear the UI). - */ - - lws_start_foreach_dll(struct lws_dll2 *, p, builder.sai_plat_owner.head) { - struct sai_plat *sp = lws_container_of(p, sai_plat_t, sai_plat_list); - sai_plat_server_ref_t *ref = NULL; - struct lwsac *ac = NULL; - sai_load_report_t lr; - char is_active = 0; - - /* - * Find the specific ref for this platform and this server - * connection (spm) - */ - lws_start_foreach_dll(struct lws_dll2 *, s, sp->servers.head) { - sai_plat_server_ref_t *r = lws_container_of(s, - sai_plat_server_ref_t, list); - if (r->spm == spm) { - ref = r; - break; - } - } lws_end_foreach_dll(s); - - if (!ref) /* This platform doesn't use this server connection */ - continue; - - /* - * Check for active tasks on this platform for this server conn - */ - lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { - struct sai_nspawn *ns = lws_container_of(d, - struct sai_nspawn, list); - if (ns->spm == spm && - ns->state == NSSTATE_EXECUTING_STEPS && ns->task) { - is_active = 1; - break; - } - } lws_end_foreach_dll(d); - - if (is_active) - any_platform_on_this_spm_active = 1; - - if (!is_active && !ref->was_active) - goto around; - - /* This platform is active for this spm, or just became idle */ - - memset(&lr, 0, sizeof(lr)); - - lws_strncpy(lr.builder_name, sp->name, sizeof(lr.builder_name)); - lr.core_count = saib_get_cpu_count(); - lr.initial_free_ram_kib = saib_get_total_ram_kib(); - lr.initial_free_disk_kib = saib_get_total_disk_kib(builder.home); - lr.reserved_ram_kib = 0; - lr.reserved_disk_kib = 0; - lr.cpu_percent = (unsigned int)saib_get_system_cpu(&builder); - lr.active_steps = 0; - lws_dll2_owner_clear(&lr.active_tasks); - - if (is_active) { - lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { - struct sai_nspawn *ns = lws_container_of(d, - struct sai_nspawn, list); - if (ns->spm == spm && - ns->state == NSSTATE_EXECUTING_STEPS && ns->task) { - sai_active_task_info_t *ati = lwsac_use_zero(&ac, sizeof(*ati), 512); - if (ati) { - lws_strncpy(ati->task_uuid, ns->task->uuid, sizeof(ati->task_uuid)); - lws_strncpy(ati->task_name, ns->task->taskname, sizeof(ati->task_name)); - ati->build_step = ns->current_step; - ati->total_steps = ns->build_step_count; - ati->est_peak_mem_kib = ns->task->est_peak_mem_kib; - ati->est_disk_kib = ns->task->est_disk_kib; - ati->started = ns->task->started; - lws_dll2_add_tail(&ati->list, &lr.active_tasks); - lr.active_steps++; - - lr.reserved_ram_kib += ns->task->est_peak_mem_kib; - lr.reserved_disk_kib += ns->task->est_disk_kib; - } - } - } lws_end_foreach_dll(d); - } - - n = saib_srv_queue_json_fragments_helper(spm->ss, - lsm_schema_json_loadreport, - LWS_ARRAY_SIZE(lsm_schema_json_loadreport), &lr); - - lwsac_free(&ac); - - if (n) - lwsl_warn("%s: failed to queue fragments\n", __func__); - - ref->was_active = is_active; - -around: - ; - } lws_end_foreach_dll(p); - - if (any_platform_on_this_spm_active) - /* Reschedule the timer only if at least one active instance */ - lws_sul_schedule(builder.context, 0, &spm->sul_load_report, - saib_sul_load_report_cb, SAI_LOAD_REPORT_US); -} - -static lws_ss_state_return_t -saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, - lws_ss_tx_ordinal_t ack) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)userobj; - struct lejp_ctx *ctx; - struct jpargs *a; - const char *pq; - int n; - - // lwsl_user("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name(state), - // (unsigned int)ack); - - switch (state) { - - case LWSSSCS_CREATING: - ctx = (struct lejp_ctx *)spm->opaque_data; - a = (struct jpargs *)ctx->user; - - /* - * Since we're "nailed up", we'll try to initiate the connection - * straight away after calling back CREATING... so we need to - * initialize any metadata etc here. - */ - - spm->index = a->next_server_index++; - - /* hook the ss up to the server url */ - - spm->url = lwsac_use(&a->builder->conf_head, - 2 *((unsigned int)ctx->npos + 1), 512); - memcpy((char *)spm->url, ctx->buf, ctx->npos); - ((char *)spm->url)[ctx->npos] = '\0'; - - lwsl_notice("%s: binding ss to %s\n", __func__, spm->url); - if (lws_ss_set_metadata(spm->ss, "url", spm->url, strlen(spm->url))) - lwsl_warn("%s: unable to set metadata\n", __func__); - - pq = spm->url; - while (*pq && (pq[0] != '/' || pq[1] != '/')) - pq++; - - if (*pq) { - n = 0; - pq += 2; - while (pq[n] && pq[n] != '/') - n++; - } else { - pq = spm->url; - n = ctx->npos; - } - - spm->name = spm->url + ctx->npos + 1; - memcpy((char *)spm->name, pq, (unsigned int)n); - ((char *)spm->name)[n] = '\0'; - - while (strchr(spm->name, '.')) - *strchr(spm->name, '.') = '_'; - while (strchr(spm->name, '/')) - *strchr(spm->name, '/') = '_'; - - /* add us to the builder list of unique servers */ - lws_dll2_add_head(&spm->list, &a->builder->sai_plat_server_owner); - - /* add us to this platforms's list of servers it accepts */ - a->mref->spm = spm; - spm->refcount++; - lws_dll2_add_tail(&a->mref->list, &a->sai_plat->servers); - - break; - - case LWSSSCS_DESTROYING: - - /* - * If the logical SS itself is going down, every platform that - * used us to connect to their server and has nspawns are also - * going down - */ - lws_dll2_foreach_safe(&builder.sai_plat_owner, spm, - cleanup_on_ss_destroy); - - break; - - case LWSSSCS_CONNECTED: - lwsl_ss_user(spm->ss, "CONNECTED"); - /* Initialize the load report SUL timer for this server connection */ - lws_sul_schedule(builder.context, 0, &spm->sul_load_report, - saib_sul_load_report_cb, 1); - - if (saib_srv_queue_json_fragments_helper(spm->ss, - lsm_schema_map_plat, - LWS_ARRAY_SIZE(lsm_schema_map_plat), - &builder.sai_plat_owner)) - return -1; - - return 0; - - case LWSSSCS_DISCONNECTED: - /* - * clean up any ongoing spawns related to this connection - */ - - lwsl_ss_user(spm->ss, "DISCONNECTED"); - lws_sul_cancel(&spm->sul_load_report); - lws_dll2_foreach_safe(&builder.sai_plat_owner, spm, - cleanup_on_ss_disconnect); - if (lws_ss_request_tx(spm->ss)) - lwsl_err("%s: failed to reconnect\n", __func__); - break; - - case LWSSSCS_ALL_RETRIES_FAILED: - lwsl_user("%s: LWSSSCS_ALL_RETRIES_FAILED\n", __func__); - return lws_ss_request_tx(spm->ss); - - case LWSSSCS_QOS_ACK_REMOTE: - lwsl_notice("%s: LWSSSCS_QOS_ACK_REMOTE\n", __func__); - break; - - default: - break; - } - - return LWSSSSRET_OK; -} - -const lws_ss_info_t ssi_sai_builder = { - .handle_offset = offsetof(struct sai_plat_server, ss), - .opaque_user_data_offset = offsetof(struct sai_plat_server, opaque_data), - .rx = saib_m_rx, - .tx = saib_m_tx, - .state = saib_m_state, - .user_alloc = sizeof(struct sai_plat_server), - .streamtype = "sai_builder" -}; diff --git a/src/builder/b-conf.c b/src/builder/b-conf.c index b343609..a887912 100644 --- a/src/builder/b-conf.c +++ b/src/builder/b-conf.c @@ -416,4 +416,8 @@ void saib_config_destroy(struct sai_builder *builder) { lwsac_free(&builder->conf_head); + +#if defined(__APPLE__) + sul_release_wakelock_cb(NULL); +#endif } diff --git a/src/builder/b-deletion.c b/src/builder/b-deletion.c index 8bfa236..dd821d9 100644 --- a/src/builder/b-deletion.c +++ b/src/builder/b-deletion.c @@ -69,7 +69,6 @@ sai_deletion_worker(const char *home_dir) * On Windows, stdin is not a pipe from the parent but a handle * value passed on the commandline */ - // detach from console... FreeConsole(); #endif @@ -152,7 +151,6 @@ int scan_jobs_dir_cb(const char *dirpath, void *user, struct lws_dir_entry *lde) { struct active_job_uuids *active = (struct active_job_uuids *)user; - struct active_job_uuid *aj; char path[512]; struct stat sb; @@ -160,10 +158,12 @@ scan_jobs_dir_cb(const char *dirpath, void *user, struct lws_dir_entry *lde) return 0; lws_start_foreach_dll(struct lws_dll2 *, p, active->owner.head) { - aj = lws_container_of(p, struct active_job_uuid, list); + struct active_job_uuid *aj = lws_container_of(p, struct active_job_uuid, list); + if (!strcmp(aj->uuid, lde->name)) /* it's an active job, leave it alone */ return 0; + } lws_end_foreach_dll(p); lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name); @@ -173,29 +173,23 @@ scan_jobs_dir_cb(const char *dirpath, void *user, struct lws_dir_entry *lde) /* older than 24h? */ if (((uint64_t)lws_now_secs() - (uint64_t)sb.st_mtime) > SAI_CLEANUP_JOB_DIR_MIN_AGE_SECS) { + char temp[128]; + size_t len = (size_t)lws_snprintf(temp, sizeof(temp), "%s\n", lde->name); +#if defined(WIN32) + DWORD written; +#endif + lwsl_notice("%s: requesting removal of old job dir %s\n", __func__, path); -#if !defined(WIN32) - { - char temp[128]; - int len = lws_snprintf(temp, sizeof(temp), "%s\n", lde->name); - if (write(builder.pipe_master_wr, temp, (unsigned int)len) != len) - lwsl_err("%s: failed to write to deletion worker\n", - __func__); - } +#if !defined(WIN32) + if (write(builder.pipe_master_wr, temp, LWS_POSIX_LENGTH_CAST(len)) != (ssize_t)len) #else - { - char temp[128]; - int len = lws_snprintf(temp, sizeof(temp), "%s\n", lde->name); - DWORD written; - - if (!WriteFile(builder.pipe_master_wr_win, temp, (DWORD)len, - &written, NULL) || written != (DWORD)len) - lwsl_err("%s: failed to write to deletion worker\n", - __func__); - } + if (!WriteFile(builder.pipe_master_wr_win, temp, (DWORD)len, + &written, NULL) || written != (DWORD)len) #endif + lwsl_err("%s: failed to write to deletion worker\n", + __func__); } return 0; diff --git a/src/builder/b-metrics.c b/src/builder/b-metrics.c index c2164b4..4ec2a3d 100644 --- a/src/builder/b-metrics.c +++ b/src/builder/b-metrics.c @@ -1,7 +1,7 @@ /* * sai-builder-metrics * - * Copyright (C) 2019 - 2021 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index b80d7cb..33e58a1 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -1,7 +1,7 @@ /* * sai-builder * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -71,9 +71,7 @@ saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel) "\"finished\":%d,", ns->retcode); n += lws_snprintf(lj + LWS_PRE + n, sizeof(lj) - LWS_PRE - (unsigned int)n, - "\"avail_slots\":%d,\"avail_mem_kib\":%u,\"avail_sto_kib\":%u,", - (int)(ns->sp->job_limit ? ns->sp->job_limit : 6u) - - ((int)ns->sp->nspawn_owner.count - 1), + "\"avail_mem_kib\":%u,\"avail_sto_kib\":%u,", saib_get_free_ram_kib(), saib_get_free_disk_kib(builder.home)); ns->retcode_set = 0; @@ -106,24 +104,9 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, uint8_t buf[600]; int ilen; - // lwsl_warn("%s: reason %d\n", __func__, reason); - switch (reason) { case LWS_CALLBACK_RAW_CLOSE_FILE: - lwsl_info("%s: stdwsi CLOSE, ns %p, lsp: %p, wsi: %p, fd: %d, stdfd: %d\n", - __func__, op ? op->ns : NULL, op ? op->lsp : NULL, - wsi, lws_get_socket_fd(wsi), lws_spawn_get_stdfd(wsi)); -/* - ilen = lws_snprintf((char *)buf, sizeof(buf), "Stdwsi %d close\n", lws_spawn_get_stdfd(wsi)); - if (ns) { - saib_log_chunk_create(ns, buf, (size_t)ilen, 3); - if (ns->spm) - if (lws_ss_request_tx(ns->spm->ss)) - lwsl_warn("%s: lws_ss_request_tx failed\n", - __func__); - } -*/ if (op && op->lsp) { if (lws_spawn_stdwsi_closed(op->lsp, wsi) && ns->reap_cb_called) { @@ -156,7 +139,8 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, len = (unsigned int)ilen; if (!op || !op->ns || !op->ns->spm) { - printf("%s: (%d) %.*s\n", __func__, (int)lws_spawn_get_stdfd(wsi), (int)len, buf); + printf("%s: (%d) %.*s\n", __func__, + (int)lws_spawn_get_stdfd(wsi), (int)len, buf); return -1; } @@ -188,6 +172,8 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, struct saib_opaque_spawn *op = (struct saib_opaque_spawn *)opaque; struct sai_nspawn *ns = op ? op->ns : NULL; uint64_t us_wallclock = op ? (uint64_t)(lws_now_usecs() - op->start_time) : 0; + char h5[40], h6[40], h7[40], h8[40], h9[40]; + sai_build_metric_t m; int exit_code = -1; char s[256]; int n; @@ -259,89 +245,60 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, peak_mem_bytes /= 1024; #endif - { - // char h1[40], h2[40], h3[40], h4[40], h10[40]; - char h5[40], h6[40], h7[40], h8[40], h9[40]; - - ns->us_cpu_user += res->us_cpu_user; - ns->us_cpu_sys += res->us_cpu_sys; - ns->us_wallclock += us_wallclock; - - if (du.size_in_bytes > ns->worst_stg) - ns->worst_stg = du.size_in_bytes; - if (peak_mem_bytes > ns->worst_mem) - ns->worst_mem = peak_mem_bytes; - - // lws_humanize_pad(h1, sizeof(h1), ns->us_cpu_user, humanize_schema_us); - // lws_humanize_pad(h2, sizeof(h2), ns->us_cpu_sys, humanize_schema_us); - // lws_humanize_pad(h3, sizeof(h3), ns->worst_mem, humanize_schema_si); - // lws_humanize_pad(h4, sizeof(h4), ns->worst_stg, humanize_schema_si); - lws_humanize_pad(h5, sizeof(h5), res->us_cpu_user, humanize_schema_us); - lws_humanize_pad(h6, sizeof(h6), res->us_cpu_sys, humanize_schema_us); - lws_humanize_pad(h7, sizeof(h7), peak_mem_bytes, humanize_schema_si); - lws_humanize_pad(h8, sizeof(h8), du.size_in_bytes, humanize_schema_si); - lws_humanize_pad(h9, sizeof(h9), us_wallclock, humanize_schema_us); - // lws_humanize_pad(h10, sizeof(h10), ns->us_wallclock, humanize_schema_us); - - n = lws_snprintf(s, sizeof(s), - ">saib> Step %d: [ %s (%s u / %s s), Mem: %sB, Stg: %sB ]\n", - ns->current_step + 1, h9, h5, h6, h7, h8); - saib_log_chunk_create(ns, s, (size_t)n, 3); - -// n = lws_snprintf(s, sizeof(s), -// ">saib> Task: [ %s (%s u / %s s), Mem: %sB, Stg: %sB ]\n", -// h10, h1, h2, h3, h4); -// saib_log_chunk_create(ns, s, (size_t)n, 3); - } + ns->us_cpu_user += res->us_cpu_user; + ns->us_cpu_sys += res->us_cpu_sys; + ns->us_wallclock += us_wallclock; - if (op->spawn) { - sai_build_metric_t m; - char hash_input[8192]; - unsigned char hash[32]; - struct lws_genhash_ctx ctx; - int n; - - if (!ns->spm) { - lwsl_err("%s: NULL ns->spm", __func__); - goto skip; - } + if (du.size_in_bytes > ns->worst_stg) + ns->worst_stg = du.size_in_bytes; + if (peak_mem_bytes > ns->worst_mem) + ns->worst_mem = peak_mem_bytes; - memset(&m, 0, sizeof(m)); - - lws_snprintf(hash_input, sizeof(hash_input), "%s%s%s%s", - ns->sp->name, op->spawn, - ns->project_name, ns->ref); - - if (lws_genhash_init(&ctx, LWS_GENHASH_TYPE_SHA256) || - lws_genhash_update(&ctx, hash_input, - strlen(hash_input)) || - lws_genhash_destroy(&ctx, hash)) - lwsl_warn("%s: sha256 failed\n", __func__); - else - for (n = 0; n < 32; n++) - lws_snprintf(m.key + (n * 2), 3, - "%02x", hash[n]); - - lws_strncpy(m.builder_name, ns->sp->name, sizeof(m.builder_name)); - lws_strncpy(m.project_name, ns->project_name, sizeof(m.project_name)); - lws_strncpy(m.ref, ns->ref, sizeof(m.ref)); - lws_strncpy(m.task_uuid, ns->task->uuid, sizeof(m.task_uuid)); - m.unixtime = (uint64_t)time(NULL); - m.us_cpu_user = res->us_cpu_user; - m.us_cpu_sys = res->us_cpu_sys; - m.wallclock_us = us_wallclock; - m.peak_mem_rss = peak_mem_bytes; - m.stg_bytes = du.size_in_bytes; - m.parallel = ns->task->parallel; - - if (saib_srv_queue_json_fragments_helper(ns->spm->ss, - lsm_schema_build_metric, - LWS_ARRAY_SIZE(lsm_schema_build_metric), &m)) - return; - } + lws_humanize_pad(h5, sizeof(h5), res->us_cpu_user, humanize_schema_us); + lws_humanize_pad(h6, sizeof(h6), res->us_cpu_sys, humanize_schema_us); + lws_humanize_pad(h7, sizeof(h7), peak_mem_bytes, humanize_schema_si); + lws_humanize_pad(h8, sizeof(h8), du.size_in_bytes, humanize_schema_si); + lws_humanize_pad(h9, sizeof(h9), us_wallclock, humanize_schema_us); + + n = lws_snprintf(s, sizeof(s), + ">saib> Step %d: [ %s (%s U / %s S), Mem: %sB, Stg: %sB ]\n", + ns->task->build_step + 1, h9, h5, h6, h7, h8); + saib_log_chunk_create(ns, s, (size_t)n, 3); + + /* + * Let's send the metrics about the step build back to the + * server so it can store them. + */ + + if (!op->spawn || !ns->spm) + goto skip; + + memset(&m, 0, sizeof(m)); + + if (sai_metrics_hash((uint8_t *)m.key, sizeof(m.key), + ns->sp->name, ns->task->build, ns->project_name, ns->ref)) + goto fail; + + lws_strncpy(m.builder_name, ns->sp->name, sizeof(m.builder_name)); + lws_strncpy(m.project_name, ns->project_name, sizeof(m.project_name)); + lws_strncpy(m.ref, ns->ref, sizeof(m.ref)); + lws_strncpy(m.task_uuid, ns->task->uuid, sizeof(m.task_uuid)); + + m.unix_time = (uint64_t)time(NULL); + m.us_cpu_user = res->us_cpu_user; + m.us_cpu_sys = res->us_cpu_sys; + m.wallclock_us = us_wallclock; + m.peak_mem_rss = peak_mem_bytes; + m.stg_bytes = du.size_in_bytes; + m.parallel = ns->task->parallel; + m.step = ns->task->build_step + 1; + + if (saib_srv_queue_json_fragments_helper(ns->spm->ss, + lsm_schema_map_build_metric, + LWS_ARRAY_SIZE(lsm_schema_build_metric), &m)) + return; skip: - ns->current_step++; /* step succeeded, wait for next instruction */ lwsl_notice("%s: step succeeded\n", __func__); @@ -375,15 +332,24 @@ skip: free(op); } - if (ns->task) + if (ns->task) { saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid, - SAI_TASK_REASON_DESTROYED); + (unsigned int)ns->retcode, + SAI_TASK_REASON_DESTROYED); + + builder.ram_reserved_kib -= ns->task->est_peak_mem_kib; + builder.disk_reserved_kib -= ns->task->est_disk_kib; + if (ns->spm) + lws_sul_schedule(builder.context, 0, + &ns->spm->sul_load_report, + saib_sul_load_report_cb, 1); + } return; fail: n = lws_snprintf(s, sizeof(s), "Build step %d FAILED, exit code: %d\n", - ns->current_step + 1, exit_code); + ns->task->build_step, exit_code); saib_log_chunk_create(ns, s, (size_t)n, 3); saib_task_grace(ns); @@ -391,9 +357,18 @@ fail: saib_log_chunk_create(ns, NULL, 0, 2); - if (ns->task) + if (ns->task) { saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid, - SAI_TASK_REASON_DESTROYED); + (unsigned int)ns->retcode, + SAI_TASK_REASON_DESTROYED); + + builder.ram_reserved_kib -= ns->task->est_peak_mem_kib; + builder.disk_reserved_kib -= ns->task->est_disk_kib; + if (ns->spm) + lws_sul_schedule(builder.context, 0, + &ns->spm->sul_load_report, + saib_sul_load_report_cb, 1); + } if (op->spawn) free(op->spawn); @@ -512,7 +487,9 @@ saib_spawn_script(struct sai_nspawn *ns) { struct lws_spawn_piped_info info; struct saib_opaque_spawn *op; - char st[2048]; +#if !defined(WIN32) + const char *script_template; +#endif const char *respath = "unk"; const char * cmd[] = { "/bin/ps", @@ -523,6 +500,8 @@ saib_spawn_script(struct sai_nspawn *ns) "LANG=en_US.UTF-8", NULL }; + char one_step[4096]; + char st[2048]; int fd, n; #if defined(__linux__) int in_cgroup = 1; @@ -539,7 +518,6 @@ saib_spawn_script(struct sai_nspawn *ns) ns->inp); #endif - char one_step[4096]; lws_strncpy(one_step, ns->task->script, sizeof(one_step)); #if defined(WIN32) @@ -564,25 +542,28 @@ saib_spawn_script(struct sai_nspawn *ns) #if defined(WIN32) n = lws_snprintf(st, sizeof(st), - ns->current_step ? runscript_win_next : runscript_win_first, + ns->task->build_step ? runscript_win_next : runscript_win_first, ns->instance_ordinal + 1, ns->task->parallel ? ns->task->parallel : 1, respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, builder.home, - ns->inp, ns->current_step > 1 ? "\\src" : "", + ns->inp, ns->task->build_step > 1 ? "\\src" : "", one_step); #else - const char *script_template; - if (ns->current_step == 0) + switch (ns->task->build_step) { + case 0: script_template = runscript_first; - else if (ns->current_step == 1) + break; + case 1: script_template = runscript_next; - else + break; + default: script_template = runscript_build; + break; + } - n = lws_snprintf(st, sizeof(st), - script_template, + n = lws_snprintf(st, sizeof(st), script_template, builder.home, ns->fsm.ovname, ns->inp_vn, ns->project_name, ns->ref, ns->instance_ordinal + 1, ns->task->parallel ? ns->task->parallel : 1, @@ -622,27 +603,24 @@ saib_spawn_script(struct sai_nspawn *ns) info.p_cgroup_ret = &in_cgroup; #endif - op = malloc(sizeof(*op)); + op = malloc(sizeof(*op)); if (!op) return 1; memset(op, 0, sizeof(*op)); - op->ns = ns; - ns->reap_cb_called = 0; - ns->op = op; + op->ns = ns; + ns->reap_cb_called = 0; + ns->op = op; #if defined(WIN32) - op->spawn = _strdup(one_step); + op->spawn = _strdup(one_step); #else - op->spawn = strdup(one_step); + op->spawn = strdup(one_step); #endif - op->start_time = lws_now_usecs(); - - info.opaque = op; - info.owner = &builder.lsp_owner; - info.plsp = &op->lsp; + op->start_time = lws_now_usecs(); - // lwsl_warn("%s: spawning build script at %llu\n", __func__, - // (unsigned long long)lws_now_usecs()); + info.opaque = op; + info.owner = &builder.lsp_owner; + info.plsp = &op->lsp; lws_spawn_piped(&info); if (!op->lsp) { @@ -651,18 +629,10 @@ saib_spawn_script(struct sai_nspawn *ns) * we can't free it here */ ns->op = NULL; - lwsl_err("%s: failed\n", __func__); return 1; } - // lwsl_warn("%s: build script spawn returned at %llu\n", __func__, - // (unsigned long long)lws_now_usecs()); - -#if defined(__linux__) - lwsl_notice("%s: lws_spawn_piped started (cgroup: %d)\n", __func__, in_cgroup); -#endif - return 0; } @@ -694,43 +664,14 @@ saib_prepare_mount(struct sai_builder *b, struct sai_nspawn *ns) goto bail_dir; lws_strncpy(homedir, ns->fsm.mp, sizeof(homedir)); - /* create work dir, session layer and mountpoint, retain mountpoint */ - -#if defined(__linux__) && 0 - - lws_snprintf(ns->fsm.mp + n, sizeof(ns->fsm.mp) - (unsigned int)n, "%cwork", csep); - m = mkdir(ns->fsm.mp, 0770); - if (m && errno != EEXIST) - goto bail_dir; - - lws_snprintf(ns->fsm.mp + n, sizeof(ns->fsm.mp) - (unsigned int)n, "%csession", csep); - m = mkdir(ns->fsm.mp, 0777); - if (m && errno != EEXIST) - goto bail_dir; - - n += lws_snprintf(ns->fsm.mp + n, sizeof(ns->fsm.mp) - (unsigned int)n, "%cmountpoint", - csep); - m = mkdir(ns->fsm.mp, 0777); - if (m && errno != EEXIST) - goto bail_dir; - - lws_snprintf(ns->fsm.mp + n, sizeof(ns->fsm.mp) - (unsigned int)n, "%chome%csai", - csep, csep); - lws_strncpy(homedir, ns->fsm.mp, sizeof(homedir)); - ns->fsm.mp[n] = '\0'; - - n = lws_fsmount_mount(&ns->fsm); - - if (!n) -#endif { /* these are ephemeral on top of the mountpoint path, we snip * them off later */ n = (int)strlen(ns->fsm.mp); - lws_snprintf(ns->fsm.mp + n, sizeof(ns->fsm.mp) - (unsigned int)n, "%chome", - csep); + lws_snprintf(ns->fsm.mp + n, sizeof(ns->fsm.mp) - (unsigned int)n, + "%chome", csep); m = mkdir(ns->fsm.mp, 0700); if (m && errno != EEXIST) goto bail_dir; diff --git a/src/builder/b-power.c b/src/builder/b-power.c index 185edd0..5238ac2 100644 --- a/src/builder/b-power.c +++ b/src/builder/b-power.c @@ -47,6 +47,8 @@ LWS_SS_USER_TYPEDEF static lws_ss_state_return_t saib_power_stay_rx(void *userobj, const uint8_t *buf, size_t len, int flags) { + char in_use = 0; + if (len < 1) return 0; @@ -74,14 +76,26 @@ saib_power_stay_rx(void *userobj, const uint8_t *buf, size_t len, int flags) struct sai_plat *sp = lws_container_of(mp, struct sai_plat, sai_plat_list); - if (sp->nspawn_owner.count) { - lwsl_warn("%s: cancelling idle grace time as ongoing task steps\n", __func__); - lws_sul_cancel(&builder.sul_idle); - return 0; + if (sp->nspawn_owner.head) { + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, sp->nspawn_owner.head) { + struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list); + + lwsl_notice("%s: ongoing task: %s\n", __func__, xns->task->uuid); + + } lws_end_foreach_dll_safe(d, d1); + + in_use = 1; } } lws_end_foreach_dll_safe(mp, mp1); + if (in_use) { + lwsl_warn("%s: cancelling idle grace time as ongoing task steps\n", __func__); + lws_sul_cancel(&builder.sul_idle); + + return 0; + } + /* * if no ongoing tasks, and we want to go OFF, then start * the idle grace timer. This will get cancelled if @@ -242,6 +256,10 @@ sul_do_suspend_cb(lws_sorted_usec_list_t *sul) void sul_idle_cb(lws_sorted_usec_list_t *sul) { +#if defined(__APPLE__) + return; +#endif + lws_ss_state_return_t r; char path[256]; @@ -338,14 +356,74 @@ saib_power_init(void) } #if defined(__APPLE__) + +int +saib_need_wakelock(void) +{ + int r = 0; + + lws_start_foreach_dll(struct lws_dll2 *, d, + builder.sai_plat_owner.head) { + struct sai_plat *sp = lws_container_of(d, struct sai_plat, sai_plat_list); + + if (sp->nspawn_owner.head) /* we are busy */ + r = 1; + + } lws_end_foreach_dll(d); + + if (builder.stay) /* there's a manual stay */ + r = 1; + + return r; +} + void sul_release_wakelock_cb(lws_sorted_usec_list_t *sul) { - lwsl_notice("%s: releasing wakelock (pid %d)\n", __func__, (int)builder.wakelock_pid); - if (builder.wakelock_pid) { - kill(builder.wakelock_pid, SIGTERM); - waitpid(builder.wakelock_pid, NULL, 0); - builder.wakelock_pid = 0; + if (!builder.wakelock_pid) + return; + + lwsl_notice("%s: releasing wakelock (pid %d)\n", __func__, + (int)builder.wakelock_pid); + + kill(builder.wakelock_pid, SIGTERM); + waitpid(builder.wakelock_pid, NULL, 0); + builder.wakelock_pid = 0; +} + +void +saib_wakelock() +{ + int need = saib_need_wakelock(); + pid_t pid; + + if (( need && builder.wakelock_pid) || + (!need && !builder.wakelock_pid)) + return; + + if (!need) { + sul_release_wakelock_cb(NULL); + return; } + + pid = fork(); + switch (pid) { + case -1: + lwsl_err("%s: fork for wakelock failed\n", __func__); + break; + case 0: + execl("/usr/bin/caffeinate", "/usr/bin/caffeinate", "-i", + (char *)NULL); + exit(1); /* should not get here */ + default: + lwsl_notice("%s: acquired wakelock (pid %d)\n", __func__, + (int)pid); + builder.wakelock_pid = pid; + break; + } + + /* if there's a pending wakelock release, cancel it */ + lws_sul_cancel(&builder.sul_release_wakelock); } + #endif diff --git a/src/builder/b-private.h b/src/builder/b-private.h index 6238489..52b0426 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -32,6 +32,7 @@ #if defined(__APPLE__) #include <sys/stat.h> /* for mkdir() */ +#include <sys/wait.h> #endif #if defined(WIN32) @@ -304,7 +305,8 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, int saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, - const char *rej_task_uuid, unsigned int reason); + const char *rej_task_uuid, unsigned int ecode, + unsigned int reason); int saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, @@ -329,3 +331,14 @@ saib_deletion_init(const char *argv0); extern void suspender_destroy(void); +#if defined(__APPLE__) +int +saib_need_wakelock(void); + +void +sul_release_wakelock_cb(lws_sorted_usec_list_t *sul); + +void +saib_wakelock(void); +#endif + diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index 5f68b9e..f48e4be 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -110,8 +110,8 @@ static const char * const default_ss_policy = "]," "\"conceal\":" "99999," "\"jitterpc\":" "20," - "\"svalidping\":" "100," - "\"svalidhup\":" "110" + "\"svalidping\":" "15," + "\"svalidhup\":" "30" "}}" "]," @@ -317,9 +317,10 @@ app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link, if (saib_deletion_init(argv0)) return 1; - +#if defined(__APPLE__) if (saib_suspender_fork(argv0)) return 1; +#endif /* * The builder JSON conf listed servers we want to connect to, @@ -370,6 +371,11 @@ app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link, lws_sul_schedule(builder.context, 0, &builder.sul_cleanup_jobs, sul_cleanup_jobs_cb, SAI_CLEANUP_JOBS_INTERVAL_US); + /* let's sample the best possible free RAM + disk situation, + * we will derate it a bit when using it */ + builder.ram_limit_kib = saib_get_free_ram_kib(); + builder.disk_total_kib = saib_get_free_disk_kib(builder.home); + break; } @@ -428,7 +434,6 @@ int main(int argc, const char **argv) if ((p = lws_cmdline_option(argc, argv, "-s"))) { lwsl_notice("%s: starting shutdown worker\n", __func__); - sleep(3000); /* * This is the suspend / shutdown worker process being spawned */ @@ -590,6 +595,10 @@ int main(int argc, const char **argv) } saib_power_init(); +#if defined(__linux__) + if (saib_suspender_fork(argv[0])) + return 1; +#endif while (!lws_service(builder.context, 0) && !interrupted) ; diff --git a/src/builder/b-suspender.c b/src/builder/b-suspender.c index bca3f30..d88cef7 100644 --- a/src/builder/b-suspender.c +++ b/src/builder/b-suspender.c @@ -63,10 +63,14 @@ saib_suspender_get_pipe(void) { #if defined(__linux__) int fd = lws_spawn_get_fd_stdxxx(lsp_suspender, 0); -#endif +#else #if defined(__APPLE__) int fd = builder.pipe_suspender_wr; +#else + int fd = 2; #endif +#endif + return fd; } @@ -121,11 +125,11 @@ callback_sai_suspender_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, struct lws_protocols protocol_suspender_stdxxx = { "sai-suspender-stdxxx", callback_sai_suspender_stdwsi, 0, 0 }; -#if !defined(__APPLE__) +#if !defined(__APPLE__) && !defined(__NetBSD__) && !defined(__OpenBSD__) static void reap(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, int we_killed_him) { - lwsl_err("%s: reaped suspender fork... %d\n", __func__, si->si_status); + // lwsl_err("%s: reaped suspender fork... %d\n", __func__, si->si_status); } #endif @@ -138,16 +142,18 @@ saib_suspender_fork(const char *path) const char * const ea[] = { rpath, "-s", NULL }; #endif +#if !defined(WIN32) if (!realpath(path, rpath)) { lwsl_err("%s: failed to get realpath for %s: %s\n", __func__, path, strerror(errno)); return 1; } +#else + lws_strncpy(rpath, path, sizeof(rpath) - 1); +#endif lwsl_err("%s: starting %s\n", __func__, rpath); - realpath(path, rpath); - #if defined(__linux__) memset(&info, 0, sizeof(info)); memset(&builder.suspend_nspawn, 0, sizeof(builder.suspend_nspawn)); diff --git a/src/builder/b-task.c b/src/builder/b-task.c index dff6a75..d5e8d94 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -26,12 +26,6 @@ #include <assert.h> #include <fcntl.h> -#if defined(__APPLE__) -#include <sys/wait.h> -void -sul_release_wakelock_cb(lws_sorted_usec_list_t *sul); -#endif - #include "b-private.h" const char *git_helper_sh = @@ -170,37 +164,32 @@ const char *git_helper_bat = static int saib_can_accept_task(sai_task_t *task, sai_plat_t *sp) { + unsigned int tc = sp->job_limit ? sp->job_limit : 6u; #if 0 unsigned int free_ram = saib_get_free_ram_kib(); unsigned int total_ram = saib_get_total_ram_kib(); unsigned int free_disk = saib_get_free_disk_kib(builder.home); unsigned int total_disk = saib_get_total_disk_kib(builder.home); // int cpu_load = saib_get_system_cpu(&builder); +#endif - if (total_ram && - (free_ram - task->est_peak_mem_kib) < (total_ram / 10) * 3) { - lwsl_notice("%s: reject task %s: not enough RAM\n", __func__, - task->uuid); - return 1; - } - - if (total_disk && - (free_disk - task->est_disk_kib) < (total_disk / 10) * 2) { - lwsl_notice("%s: reject task %s: not enough disk space\n", - __func__, task->uuid); + if ((((builder.ram_limit_kib * 4) / 3) - builder.ram_reserved_kib) < task->est_peak_mem_kib) { + lwsl_notice("%s: reject task %s: not enough RAM: task %u vs %u lim - %u res\n", __func__, + task->uuid, (unsigned int)task->est_peak_mem_kib, (unsigned int)builder.ram_limit_kib, (unsigned int)builder.ram_reserved_kib); return 1; } -/* if (cpu_load >= 0 && (cpu_load + (int)task->est_cpu_load_pct) > 50) { - lwsl_notice("%s: reject task %s: CPU load too high\n", - __func__, task->uuid); + if ((((builder.disk_total_kib * 7) / 8) - builder.disk_reserved_kib) < task->est_disk_kib) { + lwsl_notice("%s: reject task %s: not enough disk: total %u, res %u, needed %u\n", __func__, + task->uuid, (unsigned int)builder.disk_total_kib, (unsigned int)builder.disk_reserved_kib, (unsigned int)task->est_disk_kib); return 1; } -*/ -#endif - if (sp->nspawn_owner.count >= (sp->job_limit ? sp->job_limit : 6u)) + if (sp->nspawn_owner.count >= tc) { + lwsl_notice("%s: reject task %s: already running %u tasks\n", + __func__, task->uuid, tc); return 1; /* nope */ + } return 0; /* acceptable */ } @@ -252,7 +241,8 @@ saib_set_ns_state(struct sai_nspawn *ns, int state) int saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, - const char *rej_task_uuid, unsigned int reason) + const char *rej_task_uuid, unsigned int ecode, + unsigned int reason) { struct sai_rejection rej; @@ -274,11 +264,7 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, lws_snprintf(rej.host_platform, sizeof(rej.host_platform), "%s", sp->name); - rej.avail_slots = (int)(sp->job_limit ? sp->job_limit : 6u) - - (int)sp->nspawn_owner.count; - rej.avail_mem_kib = saib_get_free_ram_kib(); - rej.avail_sto_kib = saib_get_free_disk_kib(builder.home); - + rej.ecode = ecode; rej.reason = (uint8_t)reason; if (saib_srv_queue_json_fragments_helper(spm->ss, lsm_schema_json_task_rej, @@ -344,11 +330,13 @@ saib_task_destroy(struct sai_nspawn *ns) if (!m) { #if defined(__APPLE__) - lwsl_notice("%s: last task finished, scheduling wakelock release\n", __func__); - lws_sul_schedule(builder.context, 0, + if (!saib_need_wakelock()) { + lwsl_notice("%s: last task finished, scheduling wakelock release\n", __func__); + lws_sul_schedule(builder.context, 0, &builder.sul_release_wakelock, sul_release_wakelock_cb, 30 * LWS_US_PER_SEC); + } #else lws_sul_schedule(builder.context, 0, &builder.sul_idle, sul_idle_cb, @@ -405,15 +393,6 @@ saib_task_destroy(struct sai_nspawn *ns) #endif } - if (ns->task) { - builder.ram_reserved_kib -= ns->task->est_peak_mem_kib; - builder.disk_reserved_kib -= ns->task->est_disk_kib; - if (ns->spm) - lws_sul_schedule(builder.context, 0, - &ns->spm->sul_load_report, - saib_sul_load_report_cb, 1); - } - lws_dll2_remove(&ns->list); lwsl_user("%s: free(ns) %p\n", __func__, (void *)ns); free(ns); @@ -676,7 +655,6 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, sai_plat_t *sp = NULL; struct sai_nspawn *ns; int n, en, ml, fd; - sai_task_t *task; task = (sai_task_t *)a->dest; @@ -716,6 +694,7 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, lwsl_notice("%s: nspawn_census: %s\n", __func__, xns->task->uuid); } lws_end_foreach_dll_safe(d, d1); + lwsl_notice("%s:\n", __func__); /* * store a copy of the toplevel ac used for the deserialization @@ -740,17 +719,12 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { lwsl_warn("%s: server offered task that's already running\n", __func__); - saib_queue_task_status_update(sp, spm, task->uuid, + saib_queue_task_status_update(sp, spm, task->uuid, 0, SAI_TASK_REASON_DUPE); return 0; } - /* trying to reuse an nspawn? let's not do that... */ - - // if (!xns->task && !ns) - // ns = xns; - } lws_end_foreach_dll_safe(d, d1); /* @@ -759,8 +733,8 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, if (saib_can_accept_task(task, sp)) { lwsl_warn("%s: builder rejects offered task\n", __func__); - if (saib_queue_task_status_update(sp, spm, task->uuid, - SAI_TASK_REASON_BUSY)) + if (saib_queue_task_status_update(sp, spm, task->uuid, 0, + SAI_TASK_REASON_BUSY)) return -1; return 0; @@ -822,8 +796,10 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, UDS_PATHNAME_LOGPROXY"/%s.saib", #endif task->uuid); + ns->slp_control.ns = ns; ns->slp_control.log_channel_idx = 3; + if (saib_create_listen_uds(builder.context, &ns->slp_control, &ns->vhosts[0])) { lwsl_err("%s: Failed to create ctl log proxy listen UDS %s\n", @@ -853,9 +829,6 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, } } -// lwsl_hexdump_warn(task->build, strlen(task->build)); - - lws_strncpy(ns->fsm.distro, task->platform, sizeof(ns->fsm.distro)); lws_filename_purify_inplace(ns->fsm.distro); @@ -867,34 +840,28 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, * unique for remote server name ("warmcat"), * project name ("libwebsockets") */ + ns->server_name = spm->name; ns->project_name = task->repo_name; - if (!strncmp(task->git_ref, "refs/heads/", 11)) - ns->ref = task->git_ref + 11; - else - if (!strncmp(task->git_ref, "refs/tags/", 10)) - ns->ref = task->git_ref + 10; - else - ns->ref = task->git_ref; + ns->ref = sai_get_ref(task->git_ref); ns->hash = task->git_hash; ns->git_repo_url = task->git_repo_url; + if (ns->task && ns->task->ac_task_container) lwsac_free(&ns->task->ac_task_container); ns->task = task; /* we are owning this nspawn for the duration */ - ns->current_step = task->build_step; - ns->build_step_count = task->build_step_count; - if (!ns->current_step) { + ns->spm = spm; /* bind this task to the spm the req came in on */ + + if (!ns->task->build_step) { ns->spins = 0; ns->user_cancel = 0; ns->us_cpu_user = 0; ns->us_cpu_sys = 0; ns->worst_mem = 0; ns->worst_stg = 0; - } - ns->spm = spm; /* bind this task to the spm the req came in on */ - if (!ns->current_step) { + /* * If it's the first step, log some preamble info */ @@ -911,7 +878,7 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, saib_log_chunk_create(ns, ">saib>\n", 7, 3); ml = lws_snprintf(mb, sizeof(mb), ">saib> Starting task step %d ===>\n", - ns->current_step + 1); + ns->task->build_step + 1); saib_log_chunk_create(ns, mb, (unsigned int)ml, 3); @@ -1028,29 +995,16 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, * We accepted the task */ - if (saib_queue_task_status_update(sp, spm, task->uuid, SAI_TASK_REASON_ACCEPTED)) + task->started = (uint64_t)lws_now_secs(); + + builder.ram_reserved_kib += task->est_peak_mem_kib; + builder.disk_reserved_kib += task->est_disk_kib; + + if (saib_queue_task_status_update(sp, spm, task->uuid, 0, SAI_TASK_REASON_ACCEPTED)) goto bail; #if defined(__APPLE__) - /* - * If we started the first task, acquire a wakelock to prevent - * idle suspend - */ - if (sp->nspawn_owner.count == 1 && !builder.wakelock_pid) { - pid_t pid = fork(); - - if (pid == -1) - lwsl_err("%s: fork for wakelock failed\n", __func__); - else if (!pid) { - execl("/usr/bin/caffeinate", "/usr/bin/caffeinate", "-i", (char *)NULL); - exit(1); /* should not get here */ - } else { - lwsl_notice("%s: acquired wakelock (pid %d)\n", __func__, (int)pid); - builder.wakelock_pid = pid; - } - } - /* if there's a pending wakelock release, cancel it */ - lws_sul_cancel(&builder.sul_release_wakelock); + saib_wakelock(); #endif return 0; diff --git a/src/builder/b-ws-server.c b/src/builder/b-ws-server.c new file mode 100644 index 0000000..5de2db6 --- /dev/null +++ b/src/builder/b-ws-server.c @@ -0,0 +1,671 @@ +/* + * sai-builder com-warmcat-sai client protocol implementation + * + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * b1 --\ sai- sai- /-- browser + * b2 ----- server ---- web ------ browser + * b3 --/ \-- browser + * * + * + * Builder -> Server Secure Streams ws link + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> + +#include "b-private.h" + +extern struct lws_spawn_piped *lsp_suspender; + +static const lws_struct_map_t lsm_schema_json_loadreport[] = { + LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, "com.warmcat.sai.loadreport"), +}; + +static const lws_struct_map_t lsm_viewerstate_members[] = { + LSM_UNSIGNED(sai_viewer_state_t, viewers, "viewers"), +}; + +const lws_struct_map_t lsm_schema_map_m_to_b[] = { + LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), + LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, "com.warmcat.sai.taskcan"), + LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewerstate_members, + "com.warmcat.sai.viewerstate"), + LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, "com-warmcat-sai-resource"), + LSM_SCHEMA (sai_rebuild_t, NULL, lsm_rebuild, "com.warmcat.sai.rebuild") +}; + +enum { + SAIB_RX_TASK_ALLOCATION, + SAIB_RX_TASK_CANCEL, + SAIB_RX_VIEWERSTATE, + SAIB_RX_RESOURCE_REPLY, + SAIB_RX_REBUILD +}; + +/* + * This is the only path to send things from builder->server. + * + * It will copy the incoming buffer fragment into a buflist in order. So you + * should dump all your fragments for a message in here one after the other + * and the message will go out uninterrupted. Having this as the only tx path + * allows us to guarantee we won't interrupt the fragment sequencing. + * + * The fragment sizing does not have to be related to ss usage sizing, it can + * be larger and it will be used from the buflist according to what SS wants. + */ + +int +saib_srv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)lws_ss_to_user_object(h); + unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); + + *pi = ss_flags; + + // lwsl_ss_notice(h, "Queuing builder -> sai-server"); + // lwsl_hexdump_notice(buf, len); + + if (lws_buflist_append_segment(&spm->bl_to_srv, (uint8_t *)buf - sizeof(int), + len + sizeof(int)) < 0) + lwsl_ss_err(h, "failed to append"); /* still ask to drain */ + + if (lws_ss_request_tx(h)) + lwsl_ss_err(h, "failed to request tx"); + + return 0; +} + +int +saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, + const lws_struct_map_t *map, + size_t map_entries, void *object) +{ + unsigned int ssf = LWSSS_FLAG_SOM; + uint8_t buf[1024 + LWS_PRE]; + lws_struct_serialize_t *js; + size_t w = 0; + + js = lws_struct_json_serialize_create(map, map_entries, 0, object); + if (!js) { + lwsl_warn("%s: failed to serialize\n", __func__); + return -1; + } + + do { + switch (lws_struct_json_serialize(js, buf + LWS_PRE, + sizeof(buf) - LWS_PRE, &w)) { + case LSJS_RESULT_CONTINUE: + break; + case LSJS_RESULT_FINISH: + ssf |= LWSSS_FLAG_EOM; + break; + case LSJS_RESULT_ERROR: + lwsl_warn("%s: serialization failed\n", __func__); + return -1; + } + + sai_dump_stderr((const char *)buf + LWS_PRE, w); + + if (saib_srv_queue_tx(h, buf + LWS_PRE, w, ssf)) + return -1; + + ssf &= ~((unsigned int)LWSSS_FLAG_SOM); + } while (!(ssf & LWSSS_FLAG_EOM)); + + lws_struct_json_serialize_destroy(&js); + + return 0; +} + +static lws_ss_state_return_t +saib_m_rx(void *userobj, const uint8_t *in, size_t len, int flags) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)userobj; + sai_plat_t *sp = NULL; + sai_resource_t *reso; + struct lejp_ctx ctx; + lws_struct_args_t a; + sai_rebuild_t *reb; + sai_cancel_t *can; + int m; + + lws_ss_validity_confirmed(spm->ss); + + /* + * use the schema name on the incoming JSON to decide what kind of + * structure to instantiate + */ + + memset(&a, 0, sizeof(a)); + a.map_st[0] = lsm_schema_map_m_to_b; + a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_m_to_b); + a.ac_block_size = 512; + +// lwsl_hexdump_warn(in, len); + + lws_struct_json_init_parse(&ctx, NULL, &a); + m = lejp_parse(&ctx, (uint8_t *)in, (int)len); + if (m < 0) { + lwsl_hexdump_err(in, len); + lwsl_err("%s: builder rx JSON decode failed '%s'\n", + __func__, lejp_error_to_string(m)); + return m; + } + + if (!a.dest) { + lwsac_free(&a.ac); + return LWSSSSRET_OK; + } + + switch (a.top_schema_index) { + + case SAIB_RX_TASK_ALLOCATION: + if (saib_consider_allocating_task(spm, &a, in, len, flags)) + break; + + break; + + case SAIB_RX_TASK_CANCEL: + + can = (sai_cancel_t *)a.dest; + + lwsl_notice("%s: received task cancel for %s\n", __func__, can->task_uuid); + + lws_start_foreach_dll_safe(struct lws_dll2 *, mp, mp1, + builder.sai_plat_owner.head) { + struct sai_plat *sp = lws_container_of(mp, struct sai_plat, + sai_plat_list); + + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + sp->nspawn_owner.head) { + struct sai_nspawn *ns = lws_container_of(p, + struct sai_nspawn, list); + + if (ns->task && + !strcmp(can->task_uuid, ns->task->uuid)) { + lwsl_notice("%s: trying to cancel %s\n", + __func__, can->task_uuid); + + /* + * We're going to send a few signals + * at 500ms intervals + */ + ns->user_cancel = 1; + ns->term_budget = 5; + + lws_sul_schedule(ns->builder->context, 0, + &ns->sul_task_cancel, + saib_sul_task_cancel, 1); + } + + } lws_end_foreach_dll_safe(p, p1); + + } lws_end_foreach_dll_safe(mp, mp1); + break; + + case SAIB_RX_VIEWERSTATE: + { + sai_viewer_state_t *vs = (sai_viewer_state_t *)a.dest; + char any_busy = 0; + + lwsl_notice("Received viewer state update: %u viewers\n", vs->viewers); + + spm->viewer_count = vs->viewers; + + if (!vs->viewers) { + lwsl_notice("%s: VIEWERSTATE: no viewers -> no load reports\n", __func__); + lws_sul_cancel(&spm->sul_load_report); + break; + } + + /* are there any busy instances */ + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + builder.sai_plat_owner.head) { + sp = lws_container_of(d, sai_plat_t, sai_plat_list); + + lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { + struct sai_nspawn *ns = lws_container_of(d, struct sai_nspawn, list); + + if (ns->state == NSSTATE_EXECUTING_STEPS) + any_busy = 1; + + } lws_end_foreach_dll(d); + } lws_end_foreach_dll_safe(d, d1); + + if (!any_busy) { + lwsl_notice("%s: VIEWERSTATE: no busy instances -> no load reports\n", __func__); + + lws_sul_cancel(&spm->sul_load_report); + break; + } + + /* At least one viewer, start reporting */ + lwsl_notice("%s: VIEWERSTATE: viewers + busy instances -> load reports\n", __func__); + + lws_sul_schedule(builder.context, 0, &spm->sul_load_report, + saib_sul_load_report_cb, 1); + } + break; + + case SAIB_RX_RESOURCE_REPLY: + reso = (sai_resource_t *)a.dest; + + lwsl_notice("%s: RESOURCE_REPLY: cookie %s\n", + __func__, reso->cookie); + + saib_handle_resource_result(spm, (const char *)in, len); + break; + + case SAIB_RX_REBUILD: + reb = (sai_rebuild_t *)a.dest; + + lwsl_notice("%s: REBUILD: %s\n", __func__, reb->builder_name); + + if (suspender_exists) { + uint8_t b = 3; + int fd = saib_suspender_get_pipe(); + + if (write(fd, &b, 1) != 1) + lwsl_err("%s: Failed to write to suspender\n", + __func__); + } + break; + + default: + break; + } + + return LWSSSSRET_OK; +} + +/* + * We cover requested tx for any instance of a platform that can takes tasks + * from the same server... it means just by coming here, no particular + * platform / sai_plat is implied... + */ + +static lws_ss_state_return_t +saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, + int *flags) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)userobj; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&spm->bl_to_srv), depi; + char som, som1, eom, final = 1; + size_t fsl, used; + + if (!spm->bl_to_srv) + return LWSSSSRET_TX_DONT_SEND; + + depi = *pi; + *pi = (*pi) & (~(LWSSS_FLAG_SOM)); /* no SOM twice even on partial */ + + /* + * We can only issue *len at a time. + * + * Notice we are getting the stored flags from the START of the fragment each time. + * that means we can still see the right flags stored with the fragment, even if we + * have partially used the buflist frag and are partway through it. + * + * Ergo, only something to skip if we are at som=1. And also notice that although + * *pi will be right, after the lws_buflist..._use() api, what it points to has been + * destroyed. So we also dereference *pi into depi for use below. + */ + + fsl = lws_buflist_next_segment_len(&spm->bl_to_srv, NULL); + + lws_buflist_fragment_use(&spm->bl_to_srv, NULL, 0, &som, &eom); + if (som) { + fsl -= sizeof(int); + lws_buflist_fragment_use(&spm->bl_to_srv, buf, sizeof(int), &som1, &eom); + } + if (!(depi & LWSSS_FLAG_SOM)) + som = 0; + + used = (size_t)lws_buflist_fragment_use(&spm->bl_to_srv, (uint8_t *)buf, *len, &som1, &eom); + if (!used) + return LWSSSSRET_TX_DONT_SEND; + + if (used < fsl || !(depi & LWSSS_FLAG_EOM)) + final = 0; + + *len = used; + *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); + +// lwsl_ss_notice(spm->ss, "Sending %d builder->srv: ssflags %d", (int)*len, (int)*flags); +// lwsl_hexdump_notice(buf, *len); + + if (spm->bl_to_srv) + return lws_ss_request_tx(spm->ss); + + return 0; +} + +static int +cleanup_on_ss_destroy(struct lws_dll2 *d, void *user) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)user; + sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list); + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + sp->nspawn_owner.head) { + struct sai_nspawn *ns = + lws_container_of(d, struct sai_nspawn, list); + + if (ns->spm == spm) { + lwsl_warn("%s: ns->spm %p, spm %p\n", __func__, ns->spm, spm); + /* + * This pss is about to go away, make sure the ns + * can't reference it any more no matter what happens + */ + ns->spm = NULL; + } + } lws_end_foreach_dll_safe(d, d1); + + return 0; +} + +static int +cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)user; + sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list); + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, + sp->nspawn_owner.head) { + struct sai_nspawn *ns = lws_container_of(d, + struct sai_nspawn, list); + + if (ns->spm == spm) { + + /* + * This pss is about to go away, make sure the ns + * can't reference it any more no matter what happens + */ + + ns->spm = NULL; + + if (ns->op && ns->op->lsp) + lws_spawn_piped_kill_child_process(ns->op->lsp); + } + } lws_end_foreach_dll_safe(d, d1); + + return 0; +} + +void +saib_sul_load_report_cb(struct lws_sorted_usec_list *sul) +{ + struct sai_plat_server *spm = lws_container_of(sul, + struct sai_plat_server, sul_load_report); + char any_platform_on_this_spm_active = 0; + int n; + + /* + * This builder process may have multiple platforms, each with + * multiple instances. We report on each platform separately so the + * UI can distinguish them. + * + * This SUL is per-server-connection. We iterate all platforms and + * for each, see if it's supposed to connect to this server. + * + * To avoid spamming idle reports, we only report if the platform + * is active for this server, OR if it was active the last time we + * checked (ie, it has just become idle, so we need to send one last + * report with no active tasks to clear the UI). + */ + + lws_start_foreach_dll(struct lws_dll2 *, p, builder.sai_plat_owner.head) { + struct sai_plat *sp = lws_container_of(p, sai_plat_t, sai_plat_list); + sai_plat_server_ref_t *ref = NULL; + struct lwsac *ac = NULL; + sai_load_report_t lr; + char is_active = 0; + + /* + * Find the specific ref for this platform and this server + * connection (spm) + */ + lws_start_foreach_dll(struct lws_dll2 *, s, sp->servers.head) { + sai_plat_server_ref_t *r = lws_container_of(s, + sai_plat_server_ref_t, list); + if (r->spm == spm) { + ref = r; + break; + } + } lws_end_foreach_dll(s); + + if (!ref) /* This platform doesn't use this server connection */ + continue; + + /* + * Check for active tasks on this platform for this server conn + */ + lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { + struct sai_nspawn *ns = lws_container_of(d, + struct sai_nspawn, list); + if (ns->spm == spm && + ns->state == NSSTATE_EXECUTING_STEPS && ns->task) { + is_active = 1; + break; + } + } lws_end_foreach_dll(d); + + if (is_active) + any_platform_on_this_spm_active = 1; + + if (!is_active && !ref->was_active) + goto around; + + /* This platform is active for this spm, or just became idle */ + + memset(&lr, 0, sizeof(lr)); + + lws_strncpy(lr.builder_name, sp->name, sizeof(lr.builder_name)); + lr.core_count = saib_get_cpu_count(); + lr.initial_free_ram_kib = saib_get_total_ram_kib(); + lr.initial_free_disk_kib = saib_get_total_disk_kib(builder.home); + lr.reserved_ram_kib = 0; + lr.reserved_disk_kib = 0; + lr.cpu_percent = (unsigned int)saib_get_system_cpu(&builder); + lr.active_steps = 0; + lws_dll2_owner_clear(&lr.active_tasks); + + if (is_active) { + lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { + struct sai_nspawn *ns = lws_container_of(d, struct sai_nspawn, list); + + if (ns->spm == spm && + ns->state == NSSTATE_EXECUTING_STEPS && ns->task) { + sai_active_task_info_t *ati = lwsac_use_zero(&ac, sizeof(*ati), 512); + + if (ati) { + lws_strncpy(ati->task_uuid, ns->task->uuid, sizeof(ati->task_uuid)); + lws_strncpy(ati->task_name, ns->task->taskname, sizeof(ati->task_name)); + ati->build_step = ns->task->build_step; + ati->total_steps = ns->task->build_step_count; + ati->est_peak_mem_kib = ns->task->est_peak_mem_kib; + ati->est_disk_kib = ns->task->est_disk_kib; + ati->started = ns->task->started; + lws_dll2_add_tail(&ati->list, &lr.active_tasks); + lr.active_steps++; + + lr.reserved_ram_kib += ns->task->est_peak_mem_kib; + lr.reserved_disk_kib += ns->task->est_disk_kib; + } + } + } lws_end_foreach_dll(d); + } + + n = saib_srv_queue_json_fragments_helper(spm->ss, + lsm_schema_json_loadreport, + LWS_ARRAY_SIZE(lsm_schema_json_loadreport), &lr); + + lwsac_free(&ac); + + if (n) + lwsl_warn("%s: failed to queue fragments\n", __func__); + + ref->was_active = is_active; + +around: + ; + } lws_end_foreach_dll(p); + + if (any_platform_on_this_spm_active) + /* Reschedule the timer only if at least one active instance */ + lws_sul_schedule(builder.context, 0, &spm->sul_load_report, + saib_sul_load_report_cb, SAI_LOAD_REPORT_US); +} + +static lws_ss_state_return_t +saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)userobj; + struct lejp_ctx *ctx; + struct jpargs *a; + const char *pq; + int n; + + // lwsl_user("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name(state), + // (unsigned int)ack); + + switch (state) { + + case LWSSSCS_CREATING: + ctx = (struct lejp_ctx *)spm->opaque_data; + a = (struct jpargs *)ctx->user; + + /* + * Since we're "nailed up", we'll try to initiate the connection + * straight away after calling back CREATING... so we need to + * initialize any metadata etc here. + */ + + spm->index = a->next_server_index++; + + /* hook the ss up to the server url */ + + spm->url = lwsac_use(&a->builder->conf_head, + 2 *((unsigned int)ctx->npos + 1), 512); + memcpy((char *)spm->url, ctx->buf, ctx->npos); + ((char *)spm->url)[ctx->npos] = '\0'; + + lwsl_notice("%s: binding ss to %s\n", __func__, spm->url); + if (lws_ss_set_metadata(spm->ss, "url", spm->url, strlen(spm->url))) + lwsl_warn("%s: unable to set metadata\n", __func__); + + pq = spm->url; + while (*pq && (pq[0] != '/' || pq[1] != '/')) + pq++; + + if (*pq) { + n = 0; + pq += 2; + while (pq[n] && pq[n] != '/') + n++; + } else { + pq = spm->url; + n = ctx->npos; + } + + spm->name = spm->url + ctx->npos + 1; + memcpy((char *)spm->name, pq, (unsigned int)n); + ((char *)spm->name)[n] = '\0'; + + while (strchr(spm->name, '.')) + *strchr(spm->name, '.') = '_'; + while (strchr(spm->name, '/')) + *strchr(spm->name, '/') = '_'; + + /* add us to the builder list of unique servers */ + lws_dll2_add_head(&spm->list, &a->builder->sai_plat_server_owner); + + /* add us to this platforms's list of servers it accepts */ + a->mref->spm = spm; + spm->refcount++; + lws_dll2_add_tail(&a->mref->list, &a->sai_plat->servers); + + break; + + case LWSSSCS_DESTROYING: + + /* + * If the logical SS itself is going down, every platform that + * used us to connect to their server and has nspawns are also + * going down + */ + lws_dll2_foreach_safe(&builder.sai_plat_owner, spm, + cleanup_on_ss_destroy); + + break; + + case LWSSSCS_CONNECTED: + lwsl_ss_user(spm->ss, "CONNECTED"); + /* Initialize the load report SUL timer for this server connection */ + lws_sul_schedule(builder.context, 0, &spm->sul_load_report, + saib_sul_load_report_cb, 1); + + if (saib_srv_queue_json_fragments_helper(spm->ss, + lsm_schema_map_plat, + LWS_ARRAY_SIZE(lsm_schema_map_plat), + &builder.sai_plat_owner)) + return -1; + + return 0; + + case LWSSSCS_DISCONNECTED: + /* + * clean up any ongoing spawns related to this connection + */ + + lwsl_ss_user(spm->ss, "DISCONNECTED"); + lws_sul_cancel(&spm->sul_load_report); + lws_dll2_foreach_safe(&builder.sai_plat_owner, spm, + cleanup_on_ss_disconnect); + if (lws_ss_request_tx(spm->ss)) + lwsl_err("%s: failed to reconnect\n", __func__); + break; + + case LWSSSCS_ALL_RETRIES_FAILED: + lwsl_user("%s: LWSSSCS_ALL_RETRIES_FAILED\n", __func__); + return lws_ss_request_tx(spm->ss); + + case LWSSSCS_QOS_ACK_REMOTE: + lwsl_notice("%s: LWSSSCS_QOS_ACK_REMOTE\n", __func__); + break; + + default: + break; + } + + return LWSSSSRET_OK; +} + +const lws_ss_info_t ssi_sai_builder = { + .handle_offset = offsetof(struct sai_plat_server, ss), + .opaque_user_data_offset = offsetof(struct sai_plat_server, opaque_data), + .rx = saib_m_rx, + .tx = saib_m_tx, + .state = saib_m_state, + .user_alloc = sizeof(struct sai_plat_server), + .streamtype = "sai_builder" +}; diff --git a/src/common/c-utils.c b/src/common/c-utils.c index 7d815fa..671c4e5 100644 --- a/src/common/c-utils.c +++ b/src/common/c-utils.c @@ -20,8 +20,13 @@ */ #include <libwebsockets.h> + #include "include/private.h" +#if defined(WIN32) +#define write _write +#endif + int sai_uuid16_create(struct lws_context *context, char *dest33) { @@ -36,3 +41,57 @@ sai_uuid16_create(struct lws_context *context, char *dest33) return 0; } + +int +sai_metrics_hash(uint8_t *key, size_t key_len, const char *sp_name, + const char *spawn, const char *project_name, + const char *ref) +{ + struct lws_genhash_ctx ctx; + uint8_t hash[32]; + + lwsl_notice("%s: }}}}}}}}}}}}}}}}}}}}} '%s' '%s' '%s' '%s'\n", __func__, sp_name, spawn, project_name, ref); + + if (lws_genhash_init(&ctx, LWS_GENHASH_TYPE_SHA256) || + lws_genhash_update(&ctx, sp_name, strlen(sp_name)) || + lws_genhash_update(&ctx, spawn, strlen(spawn)) || + lws_genhash_update(&ctx, project_name, strlen(project_name)) || + lws_genhash_update(&ctx, ref, strlen(ref)) || + lws_genhash_destroy(&ctx, hash)) + return 1; + + lws_hex_from_byte_array(hash, sizeof(hash), (char *)key, sizeof(key_len)); + key[key_len - 1] = '\0'; + + return 0; +} + +const char * +sai_get_ref(const char *fullref) +{ + if (!strncmp(fullref, "refs/heads/", 11)) + return fullref + 11; + + if (!strncmp(fullref, "refs/tags/", 10)) + return fullref + 10; + + return fullref; +} + +const char * +sai_task_describe(sai_task_t *task, char *buf, size_t len) +{ + lws_snprintf(buf, len, "[%s(step %d/%d)]", + task->uuid, task->build_step, task->build_step_count); + + return buf; +} + +void +sai_dump_stderr(const char *buf, size_t w) +{ + if ((ssize_t)write(2, "\n", 1) != (ssize_t)1 || + (ssize_t)write(2, buf, LWS_POSIX_LENGTH_CAST(w)) != (ssize_t)w || + (ssize_t)write(2, "\n", 1) != (ssize_t)1) + lwsl_err("%s: failed to log to stderr\n", __func__); +} diff --git a/src/common/include/private.h b/src/common/include/private.h index 5181b05..f108266 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -147,8 +147,9 @@ typedef struct { /* estimations for builder resource consumption */ unsigned int est_peak_mem_kib; - unsigned int est_cpu_load_pct; unsigned int est_disk_kib; + unsigned int est_wallclock_ms; + unsigned int est_compute_ms; int parallel; char told_ongoing; @@ -212,8 +213,6 @@ struct sai_nspawn { int retcode; int instance_ordinal; int count_artifacts; - int current_step; - int build_step_count; uint8_t spins; uint8_t state; /* NSSTATE_ */ @@ -242,12 +241,12 @@ enum { typedef struct sai_rejection { struct lws_dll2 list; - char host_platform[65]; - char task_uuid[65]; - int avail_slots; - unsigned int avail_mem_kib; - unsigned int avail_sto_kib; - unsigned char reason; + char host_platform[65]; + char task_uuid[65]; + unsigned int avail_mem_kib; + unsigned int avail_sto_kib; + unsigned int ecode; + unsigned char reason; } sai_rejection_t; /* @@ -256,8 +255,8 @@ typedef struct sai_rejection { */ typedef struct sai_cancel { - struct lws_dll2 list; - char task_uuid[65]; + struct lws_dll2 list; + char task_uuid[65]; } sai_cancel_t; /* @@ -265,14 +264,14 @@ typedef struct sai_cancel { */ typedef struct sai_rebuild { - lws_dll2_t list; - char builder_name[96]; + lws_dll2_t list; + char builder_name[96]; } sai_rebuild_t; typedef struct sai_platreset { - lws_dll2_t list; - char event_uuid[65]; - char platform[65]; + lws_dll2_t list; + char event_uuid[65]; + char platform[65]; } sai_browse_rx_platreset_t; struct sai_event; @@ -304,7 +303,6 @@ typedef struct { int uid; /* builder can report this along with step completion */ - int avail_slots; unsigned int avail_mem_kib; unsigned int avail_sto_kib; } sai_log_t; @@ -485,16 +483,10 @@ typedef struct sai_plat { /* server side only: builder resource tracking */ lws_dll2_owner_t inflight_owner; /* sai_uuid_list_t */ - char last_rej_task_uuid[65]; int avail_slots; unsigned int avail_mem_kib; unsigned int avail_sto_kib; - /* server side only: for UI visibility */ - int s_avail_slots; - int s_inflight_count; - char s_last_rej_task_uuid[65]; - char windows; char power_managed; char stay_on; @@ -546,31 +538,17 @@ typedef struct sai_build_metric { char builder_name[96]; char project_name[96]; char ref[96]; - uint64_t unixtime; + uint64_t unixtime;/* actually autoincrement index in sql3 */ + uint64_t unix_time; uint64_t us_cpu_user; uint64_t us_cpu_sys; uint64_t wallclock_us; uint64_t peak_mem_rss; uint64_t stg_bytes; int parallel; + int step; } sai_build_metric_t; -typedef struct sai_build_metric_db { - lws_dll2_t list; /* for lws_struct */ - char key[65]; - char task_uuid[65]; - uint64_t unixtime; - char builder_name[96]; - char project_name[96]; - char ref[96]; - uint64_t us_cpu_user; - uint64_t us_cpu_sys; - uint64_t wallclock_us; - uint64_t peak_mem_rss; - uint64_t stg_bytes; - int parallel; -} sai_build_metric_db_t; - /* * Browser -> sai-web -> sai-server -> sai-power * @@ -620,12 +598,12 @@ extern const lws_struct_map_t lsm_schema_map_ta[1], lsm_schema_map_plat_simple[1], lsm_event[11], - lsm_task[29], - lsm_log[10], + lsm_task[30], + lsm_log[7], lsm_artifact[8], lsm_plat_list[1], lsm_schema_map_plat[1], - lsm_task_rej[6], + lsm_task_rej[4], lsm_task_cancel[1], lsm_schema_json_map_can[1], lsm_schema_json_map_task[1], @@ -635,15 +613,14 @@ extern const lws_struct_map_t lsm_rebuild[1], lsm_schema_rebuild[1], lsm_schema_build_metric[1], + lsm_schema_map_build_metric[1], lsm_schema_sq3_map_build_metric[1], lsm_load_report_members[9], lsm_schema_json_task_rej[5], lsm_stay_state_update[2], - lsm_schema_stay_state_update[1] -; -extern const lws_struct_map_t lsm_build_metric[12]; -extern const lws_struct_map_t lsm_plat[10]; -extern const lws_struct_map_t lsm_plat_for_json[16]; + lsm_schema_stay_state_update[1], + lsm_build_metric[14], + lsm_plat[13]; extern const lws_ss_info_t ssi_said_logproxy; extern struct lws_ss_handle *ssh[3]; @@ -665,4 +642,16 @@ sul_idle_cb(lws_sorted_usec_list_t *sul); int sai_uuid16_create(struct lws_context *context, char *dest33); +const char * +sai_task_describe(sai_task_t *task, char *buf, size_t len); +int +sai_metrics_hash(uint8_t *key, size_t key_len, const char *sp_name, + const char *spawn, const char *project_name, + const char *ref); + +const char * +sai_get_ref(const char *fullref); + +void +sai_dump_stderr(const char *buf, size_t w); diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 110bc1b..f5e4d38 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -1,7 +1,7 @@ /* * Sai server - ./src/common/struct-metadata.c * - * Copyright (C) 2019 - 2021 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -21,27 +21,31 @@ * lws_struct metadata for structs common to builder and server */ +#include <libwebsockets.h> + +#include "../common/include/private.h" + const lws_struct_map_t lsm_active_task_info[] = { - LSM_CARRAY (sai_active_task_info_t, task_uuid, "task_uuid"), - LSM_CARRAY (sai_active_task_info_t, task_name, "task_name"), - LSM_SIGNED (sai_active_task_info_t, build_step, "build_step"), - LSM_SIGNED (sai_active_task_info_t, total_steps, "total_steps"), - LSM_UNSIGNED (sai_active_task_info_t, est_peak_mem_kib, "est_peak_mem_kib"), - LSM_UNSIGNED (sai_active_task_info_t, est_disk_kib, "est_disk_kib"), - LSM_UNSIGNED (sai_active_task_info_t, started, "started"), + LSM_CARRAY (sai_active_task_info_t, task_uuid, "task_uuid"), + LSM_CARRAY (sai_active_task_info_t, task_name, "task_name"), + LSM_SIGNED (sai_active_task_info_t, build_step, "build_step"), + LSM_SIGNED (sai_active_task_info_t, total_steps, "total_steps"), + LSM_UNSIGNED (sai_active_task_info_t, est_peak_mem_kib, "est_peak_mem_kib"), + LSM_UNSIGNED (sai_active_task_info_t, est_disk_kib, "est_disk_kib"), + LSM_UNSIGNED (sai_active_task_info_t, started, "started"), }; const lws_struct_map_t lsm_load_report_members[] = { - LSM_CARRAY (sai_load_report_t, builder_name, "builder_name"), - LSM_SIGNED (sai_load_report_t, core_count, "core_count"), - LSM_UNSIGNED (sai_load_report_t, initial_free_ram_kib, "initial_free_ram_kib"), - LSM_UNSIGNED (sai_load_report_t, reserved_ram_kib, "reserved_ram_kib"), - LSM_UNSIGNED (sai_load_report_t, initial_free_disk_kib, "initial_free_disk_kib"), - LSM_UNSIGNED (sai_load_report_t, reserved_disk_kib, "reserved_disk_kib"), - LSM_UNSIGNED (sai_load_report_t, active_steps, "active_steps"), - LSM_UNSIGNED (sai_load_report_t, cpu_percent, "cpu_percent"), + LSM_CARRAY (sai_load_report_t, builder_name, "builder_name"), + LSM_SIGNED (sai_load_report_t, core_count, "core_count"), + LSM_UNSIGNED (sai_load_report_t, initial_free_ram_kib, "initial_free_ram_kib"), + LSM_UNSIGNED (sai_load_report_t, reserved_ram_kib, "reserved_ram_kib"), + LSM_UNSIGNED (sai_load_report_t, initial_free_disk_kib, "initial_free_disk_kib"), + LSM_UNSIGNED (sai_load_report_t, reserved_disk_kib, "reserved_disk_kib"), + LSM_UNSIGNED (sai_load_report_t, active_steps, "active_steps"), + LSM_UNSIGNED (sai_load_report_t, cpu_percent, "cpu_percent"), LSM_LIST (sai_load_report_t, active_tasks, sai_active_task_info_t, list, - NULL, lsm_active_task_info, "active_tasks"), + NULL, lsm_active_task_info, "active_tasks"), }; const lws_struct_map_t lsm_build_metric[] = { @@ -51,42 +55,33 @@ const lws_struct_map_t lsm_build_metric[] = { LSM_CARRAY (sai_build_metric_t, project_name, "project_name"), LSM_CARRAY (sai_build_metric_t, ref, "ref"), LSM_UNSIGNED (sai_build_metric_t, unixtime, "unixtime"), + LSM_UNSIGNED (sai_build_metric_t, unix_time, "unix_time"), LSM_UNSIGNED (sai_build_metric_t, us_cpu_user, "us_cpu_user"), LSM_UNSIGNED (sai_build_metric_t, us_cpu_sys, "us_cpu_sys"), LSM_UNSIGNED (sai_build_metric_t, wallclock_us, "wallclock_us"), LSM_UNSIGNED (sai_build_metric_t, peak_mem_rss, "peak_mem_rss"), LSM_UNSIGNED (sai_build_metric_t, stg_bytes, "stg_bytes"), LSM_SIGNED (sai_build_metric_t, parallel, "parallel"), + LSM_UNSIGNED (sai_build_metric_t, step, "step"), }; -const lws_struct_map_t lsm_schema_build_metric[] = { +const lws_struct_map_t lsm_schema_map_build_metric[] = { LSM_SCHEMA (sai_build_metric_t, NULL, lsm_build_metric, "com.warmcat.sai.build-metric") }; -const lws_struct_map_t lsm_sq3_build_metric[] = { - LSM_CARRAY (sai_build_metric_db_t, key, "key"), - LSM_CARRAY (sai_build_metric_db_t, task_uuid, "task_uuid"), - LSM_UNSIGNED (sai_build_metric_db_t, unixtime, "unixtime"), - LSM_CARRAY (sai_build_metric_db_t, builder_name, "builder_name"), - LSM_CARRAY (sai_build_metric_db_t, project_name, "project_name"), - LSM_CARRAY (sai_build_metric_db_t, ref, "ref"), - LSM_UNSIGNED (sai_build_metric_db_t, us_cpu_user, "us_cpu_user"), - LSM_UNSIGNED (sai_build_metric_db_t, us_cpu_sys, "us_cpu_sys"), - LSM_UNSIGNED (sai_build_metric_db_t, wallclock_us, "wallclock_us"), - LSM_UNSIGNED (sai_build_metric_db_t, peak_mem_rss, "peak_mem_rss"), - LSM_UNSIGNED (sai_build_metric_db_t, stg_bytes, "stg_bytes"), - LSM_SIGNED (sai_build_metric_db_t, parallel, "parallel"), -}; - const lws_struct_map_t lsm_schema_sq3_map_build_metric[] = { - LSM_SCHEMA_DLL2 (sai_build_metric_db_t, list, NULL, lsm_sq3_build_metric, "build_metrics"), + LSM_SCHEMA_DLL2 (sai_build_metric_t, list, NULL, lsm_build_metric, "build_metrics"), }; + const lws_struct_map_t lsm_plat[] = { /* !!! keep extern length in common/include/private.h in sync */ LSM_UNSIGNED (sai_plat_t, uid, "uid"), LSM_STRING_PTR (sai_plat_t, name, "name"), LSM_STRING_PTR (sai_plat_t, platform, "platform"), + LSM_JO_SIGNED (sai_plat_t, online, "online"), LSM_UNSIGNED (sai_plat_t, last_seen, "last_seen"), + LSM_JO_SIGNED (sai_plat_t, powering_up, "powering_up"), + LSM_JO_SIGNED (sai_plat_t, powering_down, "powering_down"), LSM_CARRAY (sai_plat_t, peer_ip, "peer_ip"), LSM_CARRAY (sai_plat_t, sai_hash, "sai_hash"), LSM_CARRAY (sai_plat_t, lws_hash, "lws_hash"), @@ -95,33 +90,13 @@ const lws_struct_map_t lsm_plat[] = { /* !!! keep extern length in common/includ LSM_UNSIGNED (sai_plat_t, stay_on, "stay_on"), }; -// This is the map for serializing to JSON -const lws_struct_map_t lsm_plat_for_json[] = { - LSM_UNSIGNED(sai_plat_t, uid, "uid"), - LSM_STRING_PTR(sai_plat_t, name, "name"), - LSM_STRING_PTR(sai_plat_t, platform,"platform"), - LSM_SIGNED(sai_plat_t, online, "online"), // MUST be present - LSM_UNSIGNED(sai_plat_t, last_seen, "last_seen"), - LSM_SIGNED(sai_plat_t, powering_up, "powering_up"), - LSM_SIGNED(sai_plat_t, powering_down, "powering_down"), - LSM_CARRAY(sai_plat_t, peer_ip, "peer_ip"), - LSM_CARRAY(sai_plat_t, sai_hash, "sai_hash"), - LSM_CARRAY(sai_plat_t, lws_hash, "lws_hash"), - LSM_UNSIGNED(sai_plat_t, windows, "windows"), - LSM_UNSIGNED(sai_plat_t, power_managed, "power_managed"), - LSM_UNSIGNED(sai_plat_t, stay_on, "stay_on"), - LSM_SIGNED(sai_plat_t, s_avail_slots, "s_avail_slots"), - LSM_SIGNED(sai_plat_t, s_inflight_count, "s_inflight_count"), - LSM_CARRAY(sai_plat_t, s_last_rej_task_uuid, "s_last_rej_task_uuid"), -}; - const lws_struct_map_t lsm_schema_map_plat_simple[] = { - LSM_SCHEMA (sai_plat_t, NULL, lsm_plat_for_json, "com-warmcat-sai-ba"), + LSM_SCHEMA (sai_plat_t, NULL, lsm_plat, "com-warmcat-sai-ba"), }; const lws_struct_map_t lsm_plat_list[] = { LSM_LIST (sai_plat_owner_t, plat_owner, sai_plat_t, - sai_plat_list, NULL, lsm_plat_for_json, "builders"), + sai_plat_list, NULL, lsm_plat, "builders"), }; const lws_struct_map_t lsm_schema_map_plat[] = { @@ -190,8 +165,9 @@ const lws_struct_map_t lsm_task[] = { LSM_SIGNED (sai_task_t, build_step, "build_step"), LSM_SIGNED (sai_task_t, build_step_count, "build_step_count"), LSM_UNSIGNED (sai_task_t, est_peak_mem_kib, "est_peak_mem_kib"), - LSM_UNSIGNED (sai_task_t, est_cpu_load_pct, "est_cpu_load_pct"), LSM_UNSIGNED (sai_task_t, est_disk_kib, "est_disk_kib"), + LSM_UNSIGNED (sai_task_t, est_wallclock_ms, "est_wallclock_ms"), + LSM_UNSIGNED (sai_task_t, est_compute_ms, "est_compute_ms"), LSM_SIGNED (sai_task_t, parallel, "parallel"), LSM_SIGNED (sai_task_t, rebuildable, "rebuildable"), }; @@ -210,9 +186,7 @@ const lws_struct_map_t lsm_schema_sq3_map_task[] = { const lws_struct_map_t lsm_task_rej[] = { LSM_CARRAY (sai_rejection_t, host_platform, "host_platform"), LSM_CARRAY (sai_rejection_t, task_uuid, "task_uuid"), - LSM_JO_SIGNED (sai_rejection_t, avail_slots, "avail_slots"), - LSM_JO_UNSIGNED (sai_rejection_t, avail_mem_kib, "avail_mem_kib"), - LSM_JO_UNSIGNED (sai_rejection_t, avail_sto_kib, "avail_sto_kib"), + LSM_JO_UNSIGNED (sai_rejection_t, ecode, "ecode"), LSM_JO_UNSIGNED (sai_rejection_t, reason, "reason"), }; @@ -259,9 +233,6 @@ const lws_struct_map_t lsm_log[] = { LSM_UNSIGNED (sai_log_t, finished, "finished"), LSM_CARRAY (sai_log_t, task_uuid, "task_uuid"), LSM_STRING_PTR (sai_log_t, log, "log"), - LSM_JO_SIGNED (sai_log_t, avail_slots, "avail_slots"), - LSM_JO_UNSIGNED (sai_log_t, avail_mem_kib, "avail_mem_kib"), - LSM_JO_UNSIGNED (sai_log_t, avail_sto_kib, "avail_sto_kib"), }; const lws_struct_map_t lsm_schema_json_map_log[] = { diff --git a/src/power/CMakeLists.txt b/src/power/CMakeLists.txt index 979852e..2842700 100644 --- a/src/power/CMakeLists.txt +++ b/src/power/CMakeLists.txt @@ -7,6 +7,7 @@ set(SRCS p-comms.c p-smartplug.c p-api.c + ../common/struct-metadata.c ) set(requirements 1) diff --git a/src/power/p-comms.c b/src/power/p-comms.c index cf5633a..8508078 100644 --- a/src/power/p-comms.c +++ b/src/power/p-comms.c @@ -27,8 +27,6 @@ #include "p-private.h" -#include "../common/struct-metadata.c" - /* Map for the "powering up" message we send to the server */ static const lws_struct_map_t lsm_schema_power_state[] = { LSM_SCHEMA(sai_power_state_t, NULL, lsm_power_state, diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index 49e50a0..df2857d 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -13,11 +13,11 @@ set(SRCS s-task.c s-task-helpers.c s-central.c - s-websrv.c + s-ws-web.c s-webops.c s-resource.c - s-metrics-db.c ../common/c-utils.c + ../common/struct-metadata.c ) set(requirements 1) diff --git a/src/server/s-central.c b/src/server/s-central.c index 6b1e48e..924960b 100644 --- a/src/server/s-central.c +++ b/src/server/s-central.c @@ -120,7 +120,7 @@ sais_central_clean_abandoned(struct vhd *vhd) lws_snprintf(s, sizeof(s), "SELECT uuid FROM tasks WHERE " "(state = %d OR state = %d) AND " - "started < %llu", + "started != 0 AND started < %llu", SAIES_PASSED_TO_BUILDER, SAIES_BEING_BUILT, (unsigned long long) (lws_now_secs() - diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 30ebbca..6ff64b2 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -37,27 +37,6 @@ #include "s-private.h" -#include "../common/struct-metadata.c" - -typedef enum { - SJS_CLONING, - SJS_ASSIGNING, - SJS_WAITING, - SJS_DONE -} sai_job_state_t; - -typedef struct sai_job { - struct lws_dll2 jobs_list; - char reponame[64]; - char ref[64]; - char head[64]; - - time_t requested; - - sai_job_state_t state; - -} sai_job_t; - extern const lws_struct_map_t lsm_schema_sq3_map_event[]; extern const lws_ss_info_t ssi_server; @@ -139,7 +118,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, if (!lws_pvo_get_str(in, "task-abandoned-timeout-mins", &num)) vhd->task_abandoned_timeout_mins = (unsigned int)atoi(num); else - vhd->task_abandoned_timeout_mins = 30; + vhd->task_abandoned_timeout_mins = 3000; if (lws_pvo_get_str(in, "database", &vhd->sqlite3_path_lhs)) { lwsl_err("%s: database pvo required\n", __func__); @@ -226,8 +205,6 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, "CREATE UNIQUE INDEX IF NOT EXISTS name_idx ON builders (name)", "create builder name index"); -// sais_mark_all_builders_offline(vhd); - lwsl_notice("%s: creating server stream\n", __func__); if (lws_ss_create(vhd->context, 0, &ssi_server, vhd, @@ -517,6 +494,8 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, ssf = (lws_is_first_fragment(wsi) ? LWSSS_FLAG_SOM : 0) | (lws_is_final_fragment(wsi) ? LWSSS_FLAG_EOM : 0); + lws_validity_confirmed(wsi); + /* * A ws client sent us something... it could be a builder or * it could be sai-power. We can tell which by the `is_power` diff --git a/src/server/s-metrics-db.c b/src/server/s-metrics-db.c deleted file mode 100644 index 75f7bcb..0000000 --- a/src/server/s-metrics-db.c +++ /dev/null @@ -1,164 +0,0 @@ -/* - * Sai server metrics db - * - * Copyright (C) 2024 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - */ - -#include <libwebsockets.h> - -#include "s-private.h" -#include "s-metrics-db.h" - -static int -sais_metrics_db_prune(struct vhd *vhd, const char *key) -{ - sqlite3_stmt *stmt; - char sql[256]; - int rc, count = 0; - - if (!vhd->pdb_metrics) - return 0; - - lws_snprintf(sql, sizeof(sql), - "SELECT COUNT(*) FROM build_metrics WHERE key = ?;"); - - rc = sqlite3_prepare_v2(vhd->pdb_metrics, sql, -1, &stmt, 0); - if (rc != SQLITE_OK) { - lwsl_err("%s: failed to prepare statement: %s\n", __func__, - sqlite3_errmsg(vhd->pdb_metrics)); - return 1; - } - - sqlite3_bind_text(stmt, 1, key, -1, SQLITE_STATIC); - - if (sqlite3_step(stmt) == SQLITE_ROW) - count = sqlite3_column_int(stmt, 0); - - sqlite3_finalize(stmt); - - if (count <= 10) - return 0; - - lws_snprintf(sql, sizeof(sql), - "DELETE FROM build_metrics WHERE key = ? AND rowid IN " - "(SELECT rowid FROM build_metrics WHERE key = ? " - "ORDER BY unixtime ASC LIMIT %d);", count - 10); - - rc = sqlite3_prepare_v2(vhd->pdb_metrics, sql, -1, &stmt, 0); - if (rc != SQLITE_OK) { - lwsl_err("%s: failed to prepare statement: %s\n", __func__, - sqlite3_errmsg(vhd->pdb_metrics)); - return 1; - } - - sqlite3_bind_text(stmt, 1, key, -1, SQLITE_STATIC); - sqlite3_bind_text(stmt, 2, key, -1, SQLITE_STATIC); - - rc = sqlite3_step(stmt); - if (rc != SQLITE_DONE) { - lwsl_err("%s: failed to delete old metrics: %s\n", __func__, - sqlite3_errmsg(vhd->pdb_metrics)); - sqlite3_finalize(stmt); - return 1; - } - - sqlite3_finalize(stmt); - - return 0; -} - -int -sais_metrics_db_init(struct vhd *vhd) -{ - char db_path[PATH_MAX]; - int rc; - - if (vhd->pdb_metrics) - return 0; - - if (!vhd->sqlite3_path_lhs) - return 0; - - lws_snprintf(db_path, sizeof(db_path), "%s-build-metrics.sqlite3", - vhd->sqlite3_path_lhs); - - rc = sqlite3_open(db_path, &vhd->pdb_metrics); - if (rc != SQLITE_OK) { - lwsl_err("%s: cannot open database %s: %s\n", __func__, - db_path, sqlite3_errmsg(vhd->pdb_metrics)); - sqlite3_close(vhd->pdb_metrics); - vhd->pdb_metrics = NULL; - return 1; - } - - if (lws_struct_sq3_create_table(vhd->pdb_metrics, - lsm_schema_sq3_map_build_metric)) { - lwsl_err("%s: failed to create build_metrics table\n", __func__); - sqlite3_close(vhd->pdb_metrics); - vhd->pdb_metrics = NULL; - return 1; - } - - return 0; -} - -void -sais_metrics_db_close(void) -{ - /* This is managed by the vhd destruction */ -} - -int -sais_metrics_db_add(struct vhd *vhd, const struct sai_build_metric *m) -{ - sai_build_metric_db_t dbm; - lws_dll2_owner_t owner; - - if (!vhd->pdb_metrics) - return 0; - - memset(&dbm, 0, sizeof(dbm)); - - lws_strncpy(dbm.key, m->key, sizeof(dbm.key)); - lws_strncpy(dbm.task_uuid, m->task_uuid, sizeof(dbm.task_uuid)); - dbm.unixtime = m->unixtime; - lws_strncpy(dbm.builder_name, m->builder_name, sizeof(dbm.builder_name)); - lws_strncpy(dbm.project_name, m->project_name, sizeof(dbm.project_name)); - lws_strncpy(dbm.ref, m->ref, sizeof(dbm.ref)); - dbm.parallel = m->parallel; - dbm.us_cpu_user = m->us_cpu_user; - dbm.us_cpu_sys = m->us_cpu_sys; - dbm.wallclock_us = m->wallclock_us; - dbm.peak_mem_rss = m->peak_mem_rss; - dbm.stg_bytes = m->stg_bytes; - - lws_dll2_owner_clear(&owner); - lws_dll2_add_tail(&dbm.list, &owner); - - if (lws_struct_sq3_serialize(vhd->pdb_metrics, - lsm_schema_sq3_map_build_metric, - &owner, 0)) { - lwsl_err("%s: failed to serialize build metric\n", __func__); - return 1; - } - - if (sais_metrics_db_prune(vhd, dbm.key)) - lwsl_warn("%s: pruning metrics failed\n", __func__); - - return 0; -} diff --git a/src/server/s-metrics-db.h b/src/server/s-metrics-db.h deleted file mode 100644 index d7f8093..0000000 --- a/src/server/s-metrics-db.h +++ /dev/null @@ -1,37 +0,0 @@ -/* - * Sai server metrics db - * - * Copyright (C) 2024 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - */ - -#if !defined(__SAI_SERVER_METRICS_DB_H__) -#define __SAI_SERVER_METRICS_DB_H__ - -struct vhd; -struct sai_build_metric; - -int -sais_metrics_db_init(struct vhd *vhd); - -void -sais_metrics_db_close(void); - -int -sais_metrics_db_add(struct vhd *vhd, const struct sai_build_metric *m); - -#endif diff --git a/src/server/s-notification.c b/src/server/s-notification.c index 04ffcc5..cb0358f 100644 --- a/src/server/s-notification.c +++ b/src/server/s-notification.c @@ -1,7 +1,7 @@ /* * Sai server - src/server/notification.c * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -268,13 +268,13 @@ sai_saifile_lejp_cb(struct lejp_ctx *ctx, char reason) */ lws_strncpy(sn->t.taskname, &ctx->path[15], sizeof(sn->t.taskname)); - sn->t.prep[0] = '\0'; - sn->t.packages[0] = '\0'; - sn->t.cmake[0] = '\0'; - sn->t.cpack[0] = '\0'; - sn->t.artifacts[0] = '\0'; - sn->t.branches[0] = '\0'; - sn->explicit_platforms[0] = '\0'; + sn->t.prep[0] = '\0'; + sn->t.packages[0] = '\0'; + sn->t.cmake[0] = '\0'; + sn->t.cpack[0] = '\0'; + sn->t.artifacts[0] = '\0'; + sn->t.branches[0] = '\0'; + sn->explicit_platforms[0] = '\0'; return 0; } @@ -578,7 +578,12 @@ next_plat: ; lws_strncpy(pss->sn.t.platform, pl->name, sizeof(pss->sn.t.platform)); - memset(&pss->sn.t.list, 0, sizeof(pss->sn.t.list)); + // pss->sn.t.server_name = ; + 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; + + lws_dll2_clear(&pss->sn.t.list); lws_dll2_owner_clear(&owner); lws_dll2_add_head(&pss->sn.t.list, &owner); diff --git a/src/server/s-power.c b/src/server/s-power.c index 30153e0..b5c48c5 100644 --- a/src/server/s-power.c +++ b/src/server/s-power.c @@ -105,15 +105,15 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, lws_start_foreach_dll(struct lws_dll2 *, p2, vhd->server.builder_owner.head) { - sai_plat_t *cb = lws_container_of(p2, sai_plat_t, sai_plat_list); - const char *dot = strchr(cb->name, '.'); + sai_plat_t *sp = lws_container_of(p2, sai_plat_t, sai_plat_list); + const char *dot = strchr(sp->name, '.'); // lwsl_notice("%s: builder entry: %s\n", __func__, cb->name); - if (dot && !(bf_set & (1 << shi)) && strlen(b->name) <= (size_t)(dot - cb->name) && - !strncmp(cb->name + (dot - cb->name) - strlen(b->name), b->name, strlen(b->name))) { - lwsl_notice("%s: ++++++++++++ Setting %s .stay_on=%d\n", __func__, cb->name, b->stay_on); - cb->stay_on = b->stay_on; + if (dot && !(bf_set & (1 << shi)) && strlen(b->name) <= (size_t)(dot - sp->name) && + !strncmp(sp->name + (dot - sp->name) - strlen(b->name), b->name, strlen(b->name))) { + lwsl_notice("%s: ++++++++++++ Setting %s .stay_on=%d\n", __func__, sp->name, b->stay_on); + sp->stay_on = b->stay_on; bf_set |= (1 << shi); } shi++; @@ -127,22 +127,22 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, } case 2: { sai_stay_state_update_t *ssu = (sai_stay_state_update_t *)a.dest; - sai_plat_t *cb; + sai_plat_t *sp; lwsl_notice("%s: Received stay_state_update for %s, stay_on=%d\n", __func__, ssu->builder_name, ssu->stay_on); lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - cb = lws_container_of(p, sai_plat_t, + sp = lws_container_of(p, sai_plat_t, sai_plat_list); - const char *dot = strchr(cb->name, '.'); + const char *dot = strchr(sp->name, '.'); - if (dot && !strncmp(cb->name, ssu->builder_name, (size_t)(dot - cb->name))) { + if (dot && !strncmp(sp->name, ssu->builder_name, (size_t)(dot - sp->name))) { lwsl_notice("%s: Updating builder %s stay_on from %d to %d\n", - __func__, cb->name, cb->stay_on, ssu->stay_on); - cb->stay_on = ssu->stay_on; + __func__, sp->name, sp->stay_on, ssu->stay_on); + sp->stay_on = ssu->stay_on; sais_list_builders(vhd); break; } @@ -155,4 +155,4 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsac_free(&a.ac); return 0; -} \ No newline at end of file +} diff --git a/src/server/s-private.h b/src/server/s-private.h index 1dba71f..a9e51e8 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -325,11 +325,14 @@ int sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char force); int -sais_set_task_state(struct vhd *vhd, const char *builder_name, - const char *builder_uuid, const char *task_uuid, sai_event_state_t state, +sais_set_task_state(struct vhd *vhd, const char *task_uuid, sai_event_state_t state, uint64_t started, uint64_t duration); int +sais_bind_task_to_builder(struct vhd *vhd, const char *builder_name, + const char *builder_uuid, const char *task_uuid); + +int sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv, lws_wsmsg_info_t *info); @@ -358,7 +361,7 @@ int sais_platforms_with_tasks_pending(struct vhd *vhd); sai_plat_t * -sais_builder_from_uuid(struct vhd *vhd, const char *hostname, const char *_file, int _line); +sais_builder_from_uuid(struct vhd *vhd, const char *hostname); sai_plat_t * sais_builder_from_host(struct vhd *vhd, const char *host); @@ -368,9 +371,6 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi); void sais_set_builder_power_state(struct vhd *vhd, const char *name, int up, int down); -void -sais_mark_all_builders_offline(struct vhd *vhd); - int sql3_get_string_cb(void *user, int cols, char **values, char **name); @@ -434,3 +434,9 @@ sais_plat_busy(sai_plat_t *sp, char set); void sais_websrv_broadcast_buflist(struct lws_ss_handle *hsrv, struct lws_buflist **bl); + +int +sais_metrics_db_init(struct vhd *vhd); + +int +sais_metrics_db_prune(struct vhd *vhd, const char *key); diff --git a/src/server/s-sai.c b/src/server/s-sai.c index 8f76a09..2668880 100644 --- a/src/server/s-sai.c +++ b/src/server/s-sai.c @@ -1,7 +1,7 @@ /* * Sai server * - * Copyright (C) 2019 - 2021 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public diff --git a/src/server/s-task-helpers.c b/src/server/s-task-helpers.c index 81cd6ad..f7b7895 100644 --- a/src/server/s-task-helpers.c +++ b/src/server/s-task-helpers.c @@ -27,60 +27,340 @@ #include "s-private.h" - void sais_get_task_metrics_estimates(struct vhd *vhd, sai_task_t *task) { - struct lws_genhash_ctx ctx; char query[256], hex[65]; sqlite3_stmt *stmt; - uint8_t hash[32]; - task->est_peak_mem_kib = 256 * 1024; /* 256MiB default */ - task->est_cpu_load_pct = 10; - task->est_disk_kib = 1024 * 1024; /* 1GiB default */ + task->est_peak_mem_kib = 0; + task->est_disk_kib = 0; + task->est_wallclock_ms = 0; /* actually total us compute */ + task->est_compute_ms = 0; - if (!vhd->pdb_metrics) + if (!vhd->pdb_metrics || !task->repo_name || !task->builder[0] || + !task->taskname[0]) return; - if (!task->repo_name || !task->platform[0] || !task->taskname[0]) + if (sai_metrics_hash((uint8_t *)hex, sizeof(hex), task->repo_name, + task->builder, task->taskname, task->git_ref)) return; - if (lws_genhash_init(&ctx, LWS_GENHASH_TYPE_SHA256) || - lws_genhash_update(&ctx, (uint8_t *)task->repo_name, strlen(task->repo_name)) || - lws_genhash_update(&ctx, (uint8_t *)task->platform, strlen(task->platform) || - lws_genhash_update(&ctx, (uint8_t *)task->taskname, strlen(task->taskname)) || - lws_genhash_destroy(&ctx, hash))) - lwsl_warn("%s: sha256 failed\n", __func__); - - lws_hex_from_byte_array(hash, sizeof(hash) - 1, hex, sizeof(hex)); - hex[64] = '\0'; - lws_snprintf(query, sizeof(query), - "SELECT AVG(peak_mem_rss), AVG(us_cpu_user), " - "AVG(stg_bytes), AVG(wallclock_us) " - "FROM build_metrics WHERE key = '%s'", - hex); + "SELECT peak_mem_rss, stg_bytes, wallclock_us, us_cpu_user, us_cpu_sys " + "FROM build_metrics " + "ORDER BY unixtime DESC " + "WHERE key = '%s' and step = %d " + "LIMIT 1", hex, task->build_step + 1); if (sqlite3_prepare_v2(vhd->pdb_metrics, query, -1, &stmt, NULL) != SQLITE_OK) return; if (sqlite3_step(stmt) == SQLITE_ROW) { - uint64_t avg_us_cpu = (uint64_t)sqlite3_column_int64(stmt, 1); - uint64_t avg_wallclock = (uint64_t)sqlite3_column_int64(stmt, 3); - if (sqlite3_column_type(stmt, 0) != SQLITE_NULL) - task->est_peak_mem_kib = (unsigned int)(sqlite3_column_int(stmt, 0) / 1024); - if (avg_wallclock) - task->est_cpu_load_pct = (unsigned int)((avg_us_cpu * 100) / avg_wallclock); + task->est_peak_mem_kib = (unsigned int)sqlite3_column_int(stmt, 0); + if (sqlite3_column_type(stmt, 1) != SQLITE_NULL) + task->est_disk_kib = (unsigned int)sqlite3_column_int(stmt, 1); if (sqlite3_column_type(stmt, 2) != SQLITE_NULL) - task->est_disk_kib = (unsigned int)(sqlite3_column_int(stmt, 2) / 1024); + task->est_wallclock_ms = (unsigned int)(sqlite3_column_int(stmt, 2) / 1000); + if (sqlite3_column_type(stmt, 3) != SQLITE_NULL && + sqlite3_column_type(stmt, 4) != SQLITE_NULL) + task->est_compute_ms = (unsigned int)((sqlite3_column_int(stmt, 3) + + sqlite3_column_int(stmt, 4)) / 1000); } sqlite3_finalize(stmt); } int +sais_bind_task_to_builder(struct vhd *vhd, const char *builder_name, + const char *builder_uuid, const char *task_uuid) +{ + char update[384], esc[96], esc1[96], esc2[96], event_uuid[33]; + struct lwsac *ac = NULL; + sai_event_t *e = NULL; + lws_dll2_owner_t o; + int n, r = 1; + + /* + * Extract the event uuid from the task uuid + */ + + sai_task_uuid_to_event_uuid(event_uuid, task_uuid); + + /* + * Look up the task's event in the event database... + */ + + lws_dll2_owner_clear(&o); + lws_sql_purify(esc1, event_uuid, sizeof(esc1)); + lws_snprintf(esc2, sizeof(esc2), " and uuid='%s'", esc1); + n = lws_struct_sq3_deserialize(vhd->server.pdb, esc2, NULL, + lsm_schema_sq3_map_event, &o, &ac, 0, 1); + if (n < 0 || !o.head) { + lwsl_err("%s: failed to get task_uuid %s\n", __func__, esc1); + goto bail; + } + + e = lws_container_of(o.head, sai_event_t, list); + + /* + * Open the event-specific database on the temporary event object + */ + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, (sqlite3 **)&e->pdb)) { + lwsl_err("%s: unable to open event-specific database\n", + __func__); + + return -1; + } + + if (builder_name) + lws_sql_purify(esc, builder_name, sizeof(esc)); + else + esc[0] = '\0'; + + if (builder_uuid) + lws_sql_purify(esc1, builder_uuid, sizeof(esc1)); + else + esc1[0] = '\0'; + lws_sql_purify(esc2, task_uuid, sizeof(esc2)); + + /* + * Update the task by uuid, in the event-specific database + */ + + lws_snprintf(update, sizeof(update), + "update tasks set builder='%s',builder_name='%s' where uuid='%s'", + esc1, esc, esc2); + + if (sqlite3_exec((sqlite3 *)e->pdb, update, NULL, NULL, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + r = 0; + +bail: + if (e) + sais_event_db_close(vhd, (sqlite3 **)&e->pdb); + lwsac_free(&ac); + + return r; +} + +int +sais_set_task_state(struct vhd *vhd, const char *task_uuid, + sai_event_state_t state, uint64_t started, uint64_t duration) +{ + char update[384], esc1[96], esc2[96], esc3[32], esc4[32], event_uuid[33]; + sai_event_state_t oes, sta, task_ostate, ostate = state; + unsigned int count = 0, count_good = 0, count_bad = 0; + uint64_t started_orig = started; + struct lwsac *ac = NULL; + sai_event_t *e = NULL; + lws_dll2_owner_t o; + int n; + + /* + * Extract the event uuid from the task uuid + */ + + sai_task_uuid_to_event_uuid(event_uuid, task_uuid); + + /* + * Look up the task's event in the event database... + */ + + lws_dll2_owner_clear(&o); + lws_sql_purify(esc1, event_uuid, sizeof(esc1)); + lws_snprintf(esc2, sizeof(esc2), " and uuid='%s'", esc1); + n = lws_struct_sq3_deserialize(vhd->server.pdb, esc2, NULL, + lsm_schema_sq3_map_event, &o, &ac, 0, 1); + if (n < 0 || !o.head) { + lwsl_err("%s: failed to get task_uuid %s\n", __func__, esc1); + goto bail; + } + + e = lws_container_of(o.head, sai_event_t, list); + oes = e->state; + + /* + * Open the event-specific database on the temporary event object + */ + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, (sqlite3 **)&e->pdb)) { + lwsl_err("%s: unable to open event-specific database\n", + __func__); + + return -1; + } + + lws_sql_purify(esc2, task_uuid, sizeof(esc2)); + + esc3[0] = esc4[0] = '\0'; + + /* + * grab the current state of it for seeing if it changed + */ + lws_snprintf(update, sizeof(update), + "select state from tasks where uuid='%s'", esc2); + if (sqlite3_exec((sqlite3 *)e->pdb, update, + sql3_get_integer_cb, &task_ostate, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + if (started) { + if (started == 1) + started = 0; + lws_snprintf(esc3, sizeof(esc3), ",started=%llu", + (unsigned long long)started); + } + if (duration) { + if (duration == 1) + duration = 0; + lws_snprintf(esc4, sizeof(esc4), ",duration=%llu", + (unsigned long long)duration); + } + + /* + * Update the task by uuid, in the event-specific database + */ + + lws_snprintf(update, sizeof(update), + "update tasks set state=%d%s%s%s where uuid='%s'", state, + esc3, esc4, state == SAIES_WAITING && started_orig == 1 ? + ",build_step=0" : "", esc2); + + if (sqlite3_exec((sqlite3 *)e->pdb, update, NULL, NULL, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + /* + * We tell interested parties about logs separately. So there's only + * something to tell about change to task state if he literally changed + * the state + */ + + if (state != task_ostate) { + + if ((state == SAIES_PASSED_TO_BUILDER || + state == SAIES_BEING_BUILT) && + !vhd->sul_activity.list.owner) + lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, + sais_activity_cb, 1 * LWS_US_PER_SEC); + + lwsl_notice("%s: seen task [%s st %d -> %d\n", __func__, + task_uuid, task_ostate, state); + + sais_taskchange(vhd->h_ss_websrv, task_uuid, state); + + if (state == SAIES_SUCCESS || state == SAIES_FAIL || + state == SAIES_CANCELLED) + lws_sul_schedule(vhd->context, 0, &vhd->sul_central, + sais_central_cb, 1); + + sais_platforms_with_tasks_pending(vhd); + + /* + * So, how many tasks for this event? + */ + + if (sqlite3_exec((sqlite3 *)e->pdb, "select count(state) from tasks", + sql3_get_integer_cb, &count, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + /* + * ... how many completed well? + */ + + if (sqlite3_exec((sqlite3 *)e->pdb, "select count(state) from tasks where state == 3", + sql3_get_integer_cb, &count_good, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + /* + * ... how many failed? + */ + + if (sqlite3_exec((sqlite3 *)e->pdb, "select count(state) from tasks where state == 4", + sql3_get_integer_cb, &count_bad, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + /* + * Decide how to set the event state based on that + */ + + lwsl_notice("%s: ev %s, task %s, state %d -> %d, count %u, good %u, bad %u, oes %d\n", + __func__, event_uuid, task_uuid, task_ostate, state, + count, count_good, count_bad, (int)oes); + + sta = SAIES_BEING_BUILT; + + if (count) { + if (count == count_good) + sta = SAIES_SUCCESS; + else + if (count == count_bad) + sta = SAIES_FAIL; + else + if (count_bad) + sta = SAIES_BEING_BUILT_HAS_FAILURES; + } + + if (sta != oes) { + lwsl_notice("%s: event state changed\n", __func__); + + /* + * Update the event + */ + + lws_sql_purify(esc1, event_uuid, sizeof(esc1)); + lws_snprintf(update, sizeof(update), + "update events set state=%d where uuid='%s'", sta, esc1); + + if (sqlite3_exec(vhd->server.pdb, update, NULL, NULL, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, update, + sqlite3_errmsg(vhd->server.pdb)); + goto bail; + } + + sais_eventchange(vhd->h_ss_websrv, event_uuid, (int)sta); + } + } + + sais_event_db_close(vhd, (sqlite3 **)&e->pdb); + lwsac_free(&ac); + + if (ostate == SAIES_STEP_SUCCESS) { + lwsl_notice("%s: sais_set_task_state() is calling sais_create_and_offer_task_step()\n", __func__); + sais_create_and_offer_task_step(vhd, task_uuid, 1); + } + + return 0; + +bail: + if (e) + sais_event_db_close(vhd, (sqlite3 **)&e->pdb); + lwsac_free(&ac); + + return 1; +} + +int sais_task_cancel(struct vhd *vhd, const char *task_uuid) { sai_cancel_t *can; @@ -124,9 +404,9 @@ sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) { char event_uuid[33], builder_name[128], esc_uuid[129], q[128]; struct pss *pss_match = NULL; - sai_plat_t *cb; sqlite3 *pdb = NULL; sai_cancel_t *can; + sai_plat_t *sp; lwsl_notice("%s: builders count %d\n", __func__, vhd->builders.count); @@ -158,14 +438,20 @@ sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) } sais_event_db_close(vhd, &pdb); - cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); - if (!cb) + /* + * This frees the sqlite task from being bound to any builder + */ + + sais_bind_task_to_builder(vhd, NULL, NULL, task_uuid); + + sp = sais_builder_from_uuid(vhd, builder_name); + if (!sp) /* Builder not connected, nothing to do */ return 0; lws_start_foreach_dll(struct lws_dll2 *, p, vhd->builders.head) { struct pss *pss = lws_container_of(p, struct pss, same); - if (pss->wsi == cb->wsi) { + if (pss->wsi == sp->wsi) { pss_match = pss; break; } @@ -201,13 +487,11 @@ sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_ sqlite3 *pdb = NULL; int ret; - lwsl_notice("%s: task reset %s\n", __func__, task_uuid); + lwsl_notice("%s: ================== task reset %s\n", __func__, task_uuid); if (!task_uuid[0]) return SAI_DB_RESULT_OK; - lwsl_notice("%s: received request to reset task %s\n", __func__, task_uuid); - sai_task_uuid_to_event_uuid(event_uuid, task_uuid); if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { @@ -245,7 +529,8 @@ sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_ sais_event_db_close(vhd, &pdb); - sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 1, 1); + /* 1,1 == reset started and duration in db for task to 0 */ + sais_set_task_state(vhd, task_uuid, SAIES_WAITING, 1, 1); sais_task_stop_on_builders(vhd, task_uuid); @@ -275,9 +560,9 @@ sai_db_result_t sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) { char esc[96], cmd[256], event_uuid[33]; + struct lwsac *ac = NULL; sqlite3 *pdb = NULL; lws_dll2_owner_t o; - struct lwsac *ac = NULL; sai_task_t *task; int ret; @@ -329,7 +614,7 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) lwsac_free(&ac); sais_event_db_close(vhd, &pdb); - sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 0, 0); + sais_set_task_state(vhd, task_uuid, SAIES_WAITING, 0, 0); sais_task_stop_on_builders(vhd, task_uuid); @@ -342,3 +627,98 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) return SAI_DB_RESULT_OK; } + +int +sais_metrics_db_prune(struct vhd *vhd, const char *key) +{ + sqlite3_stmt *stmt; + char sql[256]; + int rc, count = 0; + + if (!vhd->pdb_metrics) + return 0; + + lws_snprintf(sql, sizeof(sql), + "SELECT COUNT(*) FROM build_metrics WHERE key = ?;"); + + rc = sqlite3_prepare_v2(vhd->pdb_metrics, sql, -1, &stmt, 0); + if (rc != SQLITE_OK) { + lwsl_err("%s: failed to prepare statement: %s\n", __func__, + sqlite3_errmsg(vhd->pdb_metrics)); + return 1; + } + + sqlite3_bind_text(stmt, 1, key, -1, SQLITE_STATIC); + + if (sqlite3_step(stmt) == SQLITE_ROW) + count = sqlite3_column_int(stmt, 0); + + sqlite3_finalize(stmt); + + if (count <= 10) + return 0; + + lws_snprintf(sql, sizeof(sql), + "DELETE FROM build_metrics WHERE key = ? AND rowid IN " + "(SELECT rowid FROM build_metrics WHERE key = ? " + "ORDER BY unixtime ASC LIMIT %d);", count - 10); + + rc = sqlite3_prepare_v2(vhd->pdb_metrics, sql, -1, &stmt, 0); + if (rc != SQLITE_OK) { + lwsl_err("%s: failed to prepare statement: %s\n", __func__, + sqlite3_errmsg(vhd->pdb_metrics)); + return 1; + } + + sqlite3_bind_text(stmt, 1, key, -1, SQLITE_STATIC); + sqlite3_bind_text(stmt, 2, key, -1, SQLITE_STATIC); + + rc = sqlite3_step(stmt); + if (rc != SQLITE_DONE) { + lwsl_err("%s: failed to delete old metrics: %s\n", __func__, + sqlite3_errmsg(vhd->pdb_metrics)); + sqlite3_finalize(stmt); + return 1; + } + + sqlite3_finalize(stmt); + + return 0; +} + +int +sais_metrics_db_init(struct vhd *vhd) +{ + char db_path[PATH_MAX]; + int rc; + + if (vhd->pdb_metrics) + return 0; + + if (!vhd->sqlite3_path_lhs) + return 0; + + lws_snprintf(db_path, sizeof(db_path), "%s-build-metrics.sqlite3", + vhd->sqlite3_path_lhs); + + rc = sqlite3_open(db_path, &vhd->pdb_metrics); + if (rc != SQLITE_OK) { + lwsl_err("%s: cannot open database %s: %s\n", __func__, + db_path, sqlite3_errmsg(vhd->pdb_metrics)); + sqlite3_close(vhd->pdb_metrics); + vhd->pdb_metrics = NULL; + return 1; + } + + if (lws_struct_sq3_create_table(vhd->pdb_metrics, + lsm_schema_sq3_map_build_metric)) { + lwsl_err("%s: failed to create build_metrics table\n", __func__); + sqlite3_close(vhd->pdb_metrics); + vhd->pdb_metrics = NULL; + return 1; + } + + return 0; +} + + diff --git a/src/server/s-task.c b/src/server/s-task.c index 4842c6d..b80b6cd 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -35,235 +35,6 @@ typedef struct sai_failed_task_info { const char *taskname; } sai_failed_task_info_t; -int -sais_set_task_state(struct vhd *vhd, const char *builder_name, - const char *builder_uuid, const char *task_uuid, sai_event_state_t state, - uint64_t started, uint64_t duration) -{ - char update[384], esc[96], esc1[96], esc2[96], esc3[32], esc4[32], event_uuid[33]; - sai_event_state_t oes, sta, task_ostate, ostate = state; - unsigned int count = 0, count_good = 0, count_bad = 0; - uint64_t started_orig = started; - struct lwsac *ac = NULL; - sai_event_t *e = NULL; - lws_dll2_owner_t o; - int n; - - if (state == SAIES_STEP_SUCCESS) - state = SAIES_BEING_BUILT; - - /* - * Extract the event uuid from the task uuid - */ - - sai_task_uuid_to_event_uuid(event_uuid, task_uuid); - - /* - * Look up the task's event in the event database... - */ - - lws_dll2_owner_clear(&o); - lws_sql_purify(esc1, event_uuid, sizeof(esc1)); - lws_snprintf(esc2, sizeof(esc2), " and uuid='%s'", esc1); - n = lws_struct_sq3_deserialize(vhd->server.pdb, esc2, NULL, - lsm_schema_sq3_map_event, &o, &ac, 0, 1); - if (n < 0 || !o.head) { - lwsl_err("%s: failed to get task_uuid %s\n", __func__, esc1); - goto bail; - } - - e = lws_container_of(o.head, sai_event_t, list); - oes = e->state; - - /* - * Open the event-specific database on the temporary event object - */ - - if (sais_event_db_ensure_open(vhd, event_uuid, 0, (sqlite3 **)&e->pdb)) { - lwsl_err("%s: unable to open event-specific database\n", - __func__); - - return -1; - } - - if (builder_name) - lws_sql_purify(esc, builder_name, sizeof(esc)); - else - esc[0] = '\0'; - - if (builder_uuid) - lws_sql_purify(esc1, builder_uuid, sizeof(esc1)); - else - esc1[0] = '\0'; - lws_sql_purify(esc2, task_uuid, sizeof(esc2)); - - esc3[0] = esc4[0] = '\0'; - - /* - * grab the current state of it for seeing if it changed - */ - lws_snprintf(update, sizeof(update), - "select state from tasks where uuid='%s'", esc2); - if (sqlite3_exec((sqlite3 *)e->pdb, update, - sql3_get_integer_cb, &task_ostate, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, update, - sqlite3_errmsg(vhd->server.pdb)); - goto bail; - } - - if (started) { - if (started == 1) - started = 0; - lws_snprintf(esc3, sizeof(esc3), ",started=%llu", - (unsigned long long)started); - } - if (duration) { - if (duration == 1) - duration = 0; - lws_snprintf(esc4, sizeof(esc4), ",duration=%llu", - (unsigned long long)duration); - } - - /* - * Update the task by uuid, in the event-specific database - */ - - lws_snprintf(update, sizeof(update), - "update tasks set state=%d%s%s%s%s%s%s%s%s%s where uuid='%s'", - state, builder_uuid ? ",builder='": "", - builder_uuid ? esc1 : "", - builder_uuid ? "'" : "", - builder_name ? ",builder_name='" : "", - builder_name ? esc : "", - builder_name ? "'" : "", - esc3, esc4, state == SAIES_WAITING && started_orig == 1 ? - ",build_step=0" : "", - esc2); - - if (sqlite3_exec((sqlite3 *)e->pdb, update, NULL, NULL, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, update, - sqlite3_errmsg(vhd->server.pdb)); - goto bail; - } - - /* - * We tell interested parties about logs separately. So there's only - * something to tell about change to task state if he literally changed - * the state - */ - - if (state != task_ostate) { - - if (state == SAIES_PASSED_TO_BUILDER && - !vhd->sul_activity.list.owner) - lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, - sais_activity_cb, 1 * LWS_US_PER_SEC); - - lwsl_notice("%s: seen task %s %d -> %d\n", __func__, - task_uuid, task_ostate, state); - - sais_taskchange(vhd->h_ss_websrv, task_uuid, state); - - if (state == SAIES_SUCCESS || state == SAIES_FAIL || - state == SAIES_CANCELLED) - lws_sul_schedule(vhd->context, 0, &vhd->sul_central, - sais_central_cb, 1); - - sais_platforms_with_tasks_pending(vhd); - - /* - * So, how many tasks for this event? - */ - - if (sqlite3_exec((sqlite3 *)e->pdb, "select count(state) from tasks", - sql3_get_integer_cb, &count, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, update, - sqlite3_errmsg(vhd->server.pdb)); - goto bail; - } - - /* - * ... how many completed well? - */ - - if (sqlite3_exec((sqlite3 *)e->pdb, "select count(state) from tasks where state == 3", - sql3_get_integer_cb, &count_good, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, update, - sqlite3_errmsg(vhd->server.pdb)); - goto bail; - } - - /* - * ... how many failed? - */ - - if (sqlite3_exec((sqlite3 *)e->pdb, "select count(state) from tasks where state == 4", - sql3_get_integer_cb, &count_bad, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, update, - sqlite3_errmsg(vhd->server.pdb)); - goto bail; - } - - /* - * Decide how to set the event state based on that - */ - - lwsl_notice("%s: ev %s, task %s, state %d -> %d, count %u, good %u, bad %u, oes %d\n", - __func__, event_uuid, task_uuid, task_ostate, state, - count, count_good, count_bad, (int)oes); - - sta = SAIES_BEING_BUILT; - - if (count) { - if (count == count_good) - sta = SAIES_SUCCESS; - else - if (count == count_bad) - sta = SAIES_FAIL; - else - if (count_bad) - sta = SAIES_BEING_BUILT_HAS_FAILURES; - } - - if (sta != oes) { - lwsl_notice("%s: event state changed\n", __func__); - - /* - * Update the event - */ - - lws_sql_purify(esc1, event_uuid, sizeof(esc1)); - lws_snprintf(update, sizeof(update), - "update events set state=%d where uuid='%s'", sta, esc1); - - if (sqlite3_exec(vhd->server.pdb, update, NULL, NULL, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, update, - sqlite3_errmsg(vhd->server.pdb)); - goto bail; - } - - sais_eventchange(vhd->h_ss_websrv, event_uuid, (int)sta); - } - } - - sais_event_db_close(vhd, (sqlite3 **)&e->pdb); - lwsac_free(&ac); - - if (ostate == SAIES_STEP_SUCCESS) { - lwsl_notice("%s: sais_set_task_state() is calling sais_create_and_offer_task_step()\n", __func__); - sais_create_and_offer_task_step(vhd, task_uuid, 1); - } - - return 0; - -bail: - if (e) - sais_event_db_close(vhd, (sqlite3 **)&e->pdb); - lwsac_free(&ac); - - return 1; -} - /* * Checks if a given event db contains any tasks for a given platform */ @@ -288,7 +59,7 @@ sais_event_ran_platform(struct vhd *vhd, const char *event_uuid, sais_event_db_close(vhd, &check_pdb); - lwsl_info("%s: event %s, platform %s: count %u\n", __func__, event_uuid, + lwsl_notice("%s: event %s, platform %s: count %u\n", __func__, event_uuid, platform, count); return count > 0; @@ -433,24 +204,24 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, n = lws_struct_sq3_deserialize(vhd->server.pdb, pf, "created desc ", lsm_schema_sq3_map_event, &o, &ac, 0, 10); if (n < 0 || !o.count) { - // lwsl_notice("%s: platform %s: bail1: n %d count %d\n", __func__, platform, n, o.count); + lwsl_notice("%s: platform %s: bail1: n %d count %d\n", __func__, platform, n, o.count); goto bail; } - // lwsl_notice("%s: plat %s, toplevel results %d\n", __func__, platform, o.count); + lwsl_notice("%s: plat %s, toplevel results %d\n", __func__, platform, o.count); lws_dll2_owner_clear(&failed_tasks_owner); lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { sai_event_t *e = lws_container_of(p, sai_event_t, list); - sqlite3 *pdb = NULL, *prev_pdb = NULL; char prev_event_uuid[33] = "", checked_uuid[33] = ""; + sqlite3 *pdb = NULL, *prev_pdb = NULL; char esc_repo[96], esc_ref[96]; uint64_t last_created; int m; - // lwsl_notice("candidate event %s '%s'\n", e->uuid, esc_plat); + lwsl_notice("candidate event %s '%s'\n", e->uuid, esc_plat); if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { @@ -460,7 +231,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, */ lws_snprintf(query, sizeof(query), "select count(state) from tasks where " - "state = 0 and platform = '%s'", esc_plat); + "(state = 0 or state = 9) and platform = '%s'", esc_plat); m = sqlite3_exec(pdb, query, sql3_get_integer_cb, &pending_count, NULL); if (m != SQLITE_OK) { @@ -468,6 +239,8 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, lwsl_err("%s: query failed: %d\n", __func__, m); } + lwsl_notice("%s: %s: platform: '%s' startable tasks: %d\n", __func__, e->uuid, esc_plat, pending_count); + if (pending_count > 0) { /* there are some startable tasks on this event */ lws_sql_purify(esc_repo, e->repo_name, sizeof(esc_repo)); @@ -476,6 +249,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, do { sqlite3_stmt *sm; + int pr; prev_event_uuid[0] = '\0'; lws_snprintf(query, sizeof(query), @@ -486,8 +260,11 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, /* this is the 32-char EVENT uuid coming, not a compound (64 char) task one */ - if (sqlite3_prepare_v2(vhd->server.pdb, query, -1, &sm, NULL) != SQLITE_OK) + pr = sqlite3_prepare_v2(vhd->server.pdb, query, -1, &sm, NULL); + if (pr != SQLITE_OK) { + lwsl_warn("%s: sq3 prep returned %d instead of SQLITE_OK\n", __func__, pr); break; + } if (sqlite3_step(sm) == SQLITE_ROW) { const char *u = (const char *)sqlite3_column_text(sm, 0); if (u) { @@ -501,20 +278,26 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, lws_strncpy(prev_event_uuid, (const char *)u, sizeof(prev_event_uuid)); } last_created = (uint64_t)sqlite3_column_int64(sm, 1); - } + } else + lwsl_notice("%s: no results from event check %s %s\n", __func__, esc_repo, esc_ref); + sqlite3_finalize(sm); - if (!prev_event_uuid[0]) + if (!prev_event_uuid[0]) { + lwsl_notice("%s: breaking due to NUL prev_event_uuid\n", __func__); break; - if (!sais_event_ran_platform(vhd, prev_event_uuid, esc_plat)) + } + if (!sais_event_ran_platform(vhd, prev_event_uuid, esc_plat)) { + lwsl_notice("%s: continuing due to event_ran_platform 0\n", __func__); continue; + } lws_strncpy(checked_uuid, prev_event_uuid, sizeof(checked_uuid)); break; } while (1); if (checked_uuid[0]) { - // lwsl_notice("%s: checked_uuid %s\n", __func__, checked_uuid); + lwsl_notice("%s: checked_uuid %s\n", __func__, checked_uuid); if (!sais_event_db_ensure_open(vhd, checked_uuid, 1, &prev_pdb)) { sqlite3_stmt *sm; @@ -578,16 +361,8 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, lws_dll2_owner_t owner; lws_sql_purify(esc_taskname, fti->taskname, sizeof(esc_taskname)); - lws_snprintf(pf, sizeof(pf), " and (state == 0) and (platform == '%s') and (taskname == '%s')", + lws_snprintf(pf, sizeof(pf), " and (state == 0 or state == 9) and (platform == '%s') and (taskname == '%s')", esc_plat, esc_taskname); - if (cb->last_rej_task_uuid[0]) { - char esc_uuid[130]; - - lws_sql_purify(esc_uuid, cb->last_rej_task_uuid, - sizeof(esc_uuid)); - lws_snprintf(pf + strlen(pf), sizeof(pf) - strlen(pf), - " and (uuid != '%s')", esc_uuid); - } lwsac_free(&pss->ac_alloc_task); lws_dll2_owner_clear(&owner); @@ -609,19 +384,12 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, } } lws_end_foreach_dll(p_fail); - // lwsl_notice("%s: no priority\n", __func__); + lwsl_notice("%s: no priority\n", __func__); /* We have fallen back to doing tasks earliest-first */ - lws_snprintf(pf, sizeof(pf), " and (state = 0) and (platform = '%s')", esc_plat); - if (cb->last_rej_task_uuid[0]) { - char esc_uuid[130]; + lws_snprintf(pf, sizeof(pf), " and (state = 0 or state = 9) and (platform = '%s')", esc_plat); - lws_sql_purify(esc_uuid, cb->last_rej_task_uuid, - sizeof(esc_uuid)); - lws_snprintf(pf + strlen(pf), sizeof(pf) - strlen(pf), - " and (uuid != '%s')", esc_uuid); - } lwsac_free(&pss->ac_alloc_task); lws_dll2_owner_t owner; lws_dll2_owner_clear(&owner); @@ -630,7 +398,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, &owner, &pss->ac_alloc_task, 0, 1); // lwsl_notice("%s: deser returned %d\n", __func__, n); if (owner.count && pss->ac_alloc_task) { - // lwsl_notice("%s: orig exit\n", __func__); + lwsl_notice("%s: orig exit\n", __func__); sais_event_db_close(vhd, &pdb); lwsac_free(&ac); lwsac_free(&failed_ac); @@ -641,8 +409,8 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, return &pss->alloc_task; } - } // else - // lwsl_notice("%s: platform %s: no pending count\n", __func__, platform); + } else + lwsl_notice("%s: platform %s: no pending count\n", __func__, platform); sais_event_db_close(vhd, &pdb); @@ -653,6 +421,8 @@ bail: lwsac_free(&ac); lwsac_free(&failed_ac); + lwsl_notice("%s: leaving by bail\n", __func__); + return NULL; } @@ -804,106 +574,71 @@ bail: */ int -sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, +sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *sp, const char *platform_name) { const sai_task_t *task_template; - char original_rejected_uuid[65]; sai_task_t temp_task; - int attempts = 0; - if (cb->busy) { - lwsl_wsi_warn(pss->wsi, "::::::::::::: ABORTING task alloc due to BUSY on %s", cb->name); + if (sp->busy) { + lwsl_wsi_warn(pss->wsi, "::::::::::::: ABORTING task alloc due to BUSY on %s", sp->name); return 1; } -#if 0 - if (cb->avail_slots <= 0) { - lwsl_warn("%s: builder %s has no available slots\n", __func__, - cb->name); + /* + * Look for a task for this platform, on any event that needs building + */ + + task_template = sais_task_pending(vhd, pss, sp, platform_name); + if (!task_template) { + lwsl_notice("%s: %s: can't identify pending task\n", + __func__, sp->name); return 1; } -#endif - lws_strncpy(original_rejected_uuid, cb->last_rej_task_uuid, - sizeof(original_rejected_uuid)); - - while (attempts++ < 4) { - - /* - * Look for a task for this platform, on any event that needs building - */ - - task_template = sais_task_pending(vhd, pss, cb, platform_name); - if (!task_template) { - lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid, - sizeof(cb->last_rej_task_uuid)); - return 1; - } - - /* - * We have a candidate task, check if the builder has enough - * resources for it - */ - memcpy(&temp_task, task_template, sizeof(temp_task)); - sais_get_task_metrics_estimates(vhd, &temp_task); - - if (temp_task.est_peak_mem_kib > cb->avail_mem_kib || - temp_task.est_disk_kib > cb->avail_sto_kib) { - lwsl_notice("%s: builder %s lacks resources for task %s " - "(mem %uk/%uk, sto %uk/%uk), trying another\n", - __func__, cb->name, temp_task.uuid, - temp_task.est_peak_mem_kib, cb->avail_mem_kib, - temp_task.est_disk_kib, cb->avail_sto_kib); - - /* mark it rejected for this builder and try again */ - lws_strncpy(cb->last_rej_task_uuid, temp_task.uuid, - sizeof(cb->last_rej_task_uuid)); - continue; - } - - if (sais_is_task_inflight(vhd, NULL, task_template->uuid, NULL)) { - lwsl_notice("%s: ~~~~~~~~ skipping %s as listed on inflight\n", __func__, task_template->uuid); - continue; - } - lwsl_notice("%s: %s: task %s found for %s\n", __func__, - platform_name, task_template->uuid, cb->name); - - /* yes, we will offer it to him */ - - if (sais_set_task_state(vhd, cb->name, cb->name, task_template->uuid, - SAIES_PASSED_TO_BUILDER, lws_now_secs(), 0)) - goto bail; + /* + * We have a candidate task, check if the builder has enough + * resources for it + */ + memcpy(&temp_task, task_template, sizeof(temp_task)); + sais_get_task_metrics_estimates(vhd, &temp_task); + + if (temp_task.est_peak_mem_kib > sp->avail_mem_kib || + temp_task.est_disk_kib > sp->avail_sto_kib) { + lwsl_notice("%s: builder %s lacks resources for task %s " + "(mem %uk/%uk, sto %uk/%uk), trying another\n", + __func__, sp->name, temp_task.uuid, + temp_task.est_peak_mem_kib, sp->avail_mem_kib, + temp_task.est_disk_kib, sp->avail_sto_kib); + return 1; + } - cb->s_avail_slots = cb->avail_slots; - cb->s_inflight_count = (int)cb->inflight_owner.count; - lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, - sizeof(cb->s_last_rej_task_uuid)); + if (sais_is_task_inflight(vhd, NULL, task_template->uuid, NULL)) { + lwsl_notice("%s: ~~~~~~~~ skipping %s as listed on inflight\n", + __func__, task_template->uuid); + return 1; + } - sais_list_builders(vhd); + /* + * This marks the sqlite task as being bound to builder sp->name + */ - /* advance the task state first time we get logs */ - pss->mark_started = 1; + sais_bind_task_to_builder(vhd, sp->name, sp->name, task_template->uuid); - sais_create_and_offer_task_step(vhd, task_template->uuid, 3); + lwsl_notice("%s: %s: task %s found for %s\n", __func__, + platform_name, task_template->uuid, sp->name); - return 0; - } + if (sais_create_and_offer_task_step(vhd, task_template->uuid, 3)) + return 1; - lwsl_warn("%s: exceeded max attempts to find suitable task for %s\n", - __func__, cb->name); - lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid, - sizeof(cb->last_rej_task_uuid)); + /* yes, we will offer it to him */ - return 1; + sais_list_builders(vhd); -bail: - lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid, - sizeof(cb->last_rej_task_uuid)); - lwsac_free(&pss->a.ac); - lwsac_free(&pss->ac_alloc_task); + /* advance the task state first time we get logs */ + pss->mark_started = 1; - return -1; + return 0; } #define MAX_BLOB 1024 @@ -1020,8 +755,9 @@ nope: int sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char force) { - char event_uuid[33], esc_uuid[129], *p, *start, url[128], mirror_path[256], update[128]; - sai_task_t *temp_task = NULL; + char event_uuid[33], esc_uuid[129], *p, *start, url[128], + mirror_path[256], update[128]; + sai_task_t *temp_task = NULL, *task_template; lws_dll2_owner_t o, o_event; struct lwsac *ac = NULL; sai_uuid_list_t *ul; @@ -1029,7 +765,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for sai_event_t *event; int n, build_step; struct pss *pss; - sai_plat_t *cb; + sai_plat_t *sp; int inflight; int ret = -1; @@ -1055,30 +791,30 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for n = lws_struct_sq3_deserialize(pdb, update, NULL, lsm_schema_sq3_map_task, &o, &ac, 0, 1); if (n < 0 || !o.head) { + lwsl_warn("%s: bailing as nothing with state != 4\n", __func__); sais_event_db_close(vhd, &pdb); lwsac_free(&ac); return -1; } - { - sai_task_t *task_template = lws_container_of(o.head, sai_task_t, list); + task_template = lws_container_of(o.head, sai_task_t, list); - /* - * Make a copy of the lws_struct allocation in the lwsac, - * then drop the lwsac - */ + /* + * Make a copy of the lws_struct allocation in the lwsac, + * then drop the lwsac + */ - temp_task = malloc(sizeof(sai_task_t)); - if (!temp_task) { - lwsac_free(&ac); - sais_event_db_close(vhd, &pdb); - return -1; - } - memset(temp_task, 0, sizeof(*temp_task)); - *temp_task = *task_template; + temp_task = malloc(sizeof(sai_task_t)); + if (!temp_task) { lwsac_free(&ac); + sais_event_db_close(vhd, &pdb); + return -1; } + memset(temp_task, 0, sizeof(*temp_task)); + *temp_task = *task_template; + lwsac_free(&ac); + sais_get_task_metrics_estimates(vhd, temp_task); build_step = temp_task->build_step; @@ -1090,33 +826,37 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for lsm_schema_sq3_map_event, &o_event, &temp_task->ac_task_container, 0, 1); if (n < 0 || !o_event.head) { + lwsl_warn("%s: bailing as nothing with uuid %s\n", __func__, esc_uuid); sais_event_db_close(vhd, &pdb); free(temp_task); return -1; } event = lws_container_of(o_event.head, sai_event_t, list); - temp_task->one_event = event; + temp_task->one_event = event; temp_task->repo_name = event->repo_name; temp_task->git_ref = event->ref; temp_task->git_hash = event->hash; temp_task->git_repo_url = event->repo_fetchurl; /* find builder */ - cb = sais_builder_from_uuid(vhd, temp_task->builder_name, __FILE__, __LINE__); - if (!cb) + sp = sais_builder_from_uuid(vhd, temp_task->builder_name); + if (!sp) { + lwsl_warn("%s: bailing as can't find builder from %s\n", __func__, temp_task->builder_name); goto bail; + } - if (sais_add_to_inflight_list_if_absent(vhd, cb, task_uuid)) { - sais_task_clear_build_and_logs(vhd, task_uuid, 0); + if (sais_is_task_inflight(vhd, NULL, task_uuid, &ul)) { + lwsl_warn("%s: bailing as inflight %s\n", __func__, task_uuid); goto bail; } - /* provisionally decrement until we hear from builder */ - if (cb->avail_slots > 0) - cb->avail_slots--; - + if (sais_add_to_inflight_list_if_absent(vhd, sp, task_uuid)) { + lwsl_warn("%s: bailing as can't add to inflight %s\n", __func__, task_uuid); + sais_task_clear_build_and_logs(vhd, task_uuid, 0); + goto bail; + } lws_strncpy(url, temp_task->one_event->repo_fetchurl, sizeof(url)); lws_filename_purify_inplace(url); @@ -1130,7 +870,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for switch (build_step) { case 0: /* git mirror */ - if (cb->windows) + if (sp->windows) lws_snprintf(temp_task->script, sizeof(temp_task->script), ".\\git_helper.bat mirror \"%s\" %s %s %s", temp_task->git_repo_url, temp_task->git_ref, temp_task->git_hash, @@ -1142,7 +882,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for mirror_path); break; case 1: /* git checkout */ - if (cb->windows) + if (sp->windows) lws_snprintf(temp_task->script, sizeof(temp_task->script), ".\\git_helper.bat checkout \"%s\" src %s", mirror_path, temp_task->git_hash); @@ -1164,9 +904,9 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for lwsl_err("%s: +++++++++++++++++++ determined no more steps after build_step %d for task %s, setting SAIES_SUCCESS\n", __func__, build_step, temp_task->uuid); - sais_set_task_state(vhd, NULL, NULL, temp_task->uuid, SAIES_SUCCESS, 0, 0); + sais_set_task_state(vhd, temp_task->uuid, SAIES_SUCCESS, 0, 0); - if (sais_is_task_inflight(vhd, cb, temp_task->uuid, &u)) + if (sais_is_task_inflight(vhd, sp, temp_task->uuid, &u)) sais_inflight_entry_destroy(u); ret = 0; goto bail; @@ -1186,7 +926,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for pss = NULL; lws_start_foreach_dll(struct lws_dll2 *, d, vhd->builders.head) { struct pss *pss_ = lws_container_of(d, struct pss, same); - if (pss_->wsi == cb->wsi) { + if (pss_->wsi == sp->wsi) { pss = pss_; break; } @@ -1197,7 +937,8 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for temp_task->server_name = pss->server_name; - if (sais_add_to_inflight_list_if_absent(vhd, cb, temp_task->uuid)) { + if (sais_add_to_inflight_list_if_absent(vhd, sp, temp_task->uuid)) { + lwsl_warn("%s: bailing as can't add to inflight %s\n", __func__, task_uuid); sais_task_clear_build_and_logs(vhd, temp_task->uuid, 0); goto bail; } @@ -1209,15 +950,6 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for lws_dll2_add_tail(&temp_task->pending_assign_list, &pss->issue_task_owner); lws_callback_on_writable(pss->wsi); - /* - * If the task hasn't failed, bump the build step - */ - - lws_sql_purify(esc_uuid, task_uuid, sizeof(esc_uuid)); - lws_snprintf(update, sizeof(update), "update tasks set build_step=%d where state != 4 and uuid='%s'", - build_step + 1, esc_uuid); - sqlite3_exec(pdb, update, NULL, NULL, NULL); - sais_event_db_close(vhd, &pdb); return 0; @@ -1237,18 +969,21 @@ sais_plat_find_jobs_cb(lws_sorted_usec_list_t *sul) { sai_plat_t *sp = lws_container_of(sul, sai_plat_t, sul_find_jobs); - if (!sp->busy && sp->wsi && lws_wsi_user(sp->wsi)) + lwsl_notice("%s: %s: sp->busy: %d\n", __func__, sp->name, sp->busy); + + if (!sp->busy && sp->wsi && lws_wsi_user(sp->wsi) && /* * try to bind outstanding task to specific builder * instance */ - sais_allocate_task((struct vhd *)sp->vhd, - (struct pss *)lws_wsi_user(sp->wsi), - sp, sp->platform); - - if (!sp->busy) + !sais_allocate_task((struct vhd *)sp->vhd, + (struct pss *)lws_wsi_user(sp->wsi), + sp, sp->platform)) + /* + * Only look again if we ended this try successfully + */ lws_sul_schedule(sp->cx, 0, &sp->sul_find_jobs, - sais_plat_find_jobs_cb, 500 * LWS_US_PER_MS); + sais_plat_find_jobs_cb, 50 * LWS_US_PER_MS); } void @@ -1265,5 +1000,5 @@ sais_plat_busy(sai_plat_t *sp, char set) lwsl_notice("%s: %s: CLEARING BUSY\n", __func__, sp->name); lws_sul_schedule(sp->cx, 0, &sp->sul_find_jobs, - sais_plat_find_jobs_cb, 500 * LWS_US_PER_MS); + sais_plat_find_jobs_cb, 1); } diff --git a/src/server/s-webops.c b/src/server/s-webops.c index 82bc57c..87ba5e5 100644 --- a/src/server/s-webops.c +++ b/src/server/s-webops.c @@ -61,7 +61,7 @@ _sais_websrv_broadcast(struct lws_ss_handle *h, void *arg) 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\n", + // 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; @@ -135,24 +135,19 @@ sais_websrv_broadcast_buflist(struct lws_ss_handle *hsrv, struct lws_buflist **b } - -struct sais_arg { - const char *uid; - int state; -}; - -static void -_sais_taskchange(struct lws_ss_handle *h, void *_arg) +void +sais_taskchange(struct lws_ss_handle *hsrv, const char *task_uuid, int state) { - struct sais_arg *arg = (struct sais_arg *)_arg; - char tc[LWS_PRE + 128], *start = tc + LWS_PRE; + char tc[LWS_PRE + 256], *start = tc + LWS_PRE; lws_wsmsg_info_t info; int n; + lwsl_ss_notice(hsrv, "%%%%%%%% sai-taskchange %s -> %d", task_uuid, state); + n = lws_snprintf(start, sizeof(tc) - LWS_PRE, "{\"schema\":\"sai-taskchange\", " "\"event_hash\":\"%s\", \"state\":%d}", - arg->uid, arg->state); + task_uuid, state); memset(&info, 0, sizeof(info)); info.private_source_idx = SAI_WEBSRV_PB__GENERATED; @@ -160,36 +155,26 @@ _sais_taskchange(struct lws_ss_handle *h, void *_arg) info.len = (size_t)n; info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, &info) < 0) { + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(hsrv, &info) < 0) { lwsl_warn("%s: buflist append failed\n", __func__); return; } - - if (lws_ss_request_tx(h)) - lwsl_ss_warn(h, "tx req fail"); } void -sais_taskchange(struct lws_ss_handle *hsrv, const char *task_uuid, int state) -{ - struct sais_arg arg = { task_uuid, state }; - - lws_ss_server_foreach_client(hsrv, _sais_taskchange, (void *)&arg); -} - -static void -_sais_eventchange(struct lws_ss_handle *h, void *_arg) +sais_eventchange(struct lws_ss_handle *hsrv, const char *event_uuid, int state) { - struct sais_arg *arg = (struct sais_arg *)_arg; - char tc[LWS_PRE + 128], *start = tc + LWS_PRE; + char tc[LWS_PRE + 256], *start = tc + LWS_PRE; lws_wsmsg_info_t info; int n; + lwsl_ss_notice(hsrv, "%%%%%%%% sai-eventchange %s -> %d", event_uuid, state); + n = lws_snprintf(start, sizeof(tc) - LWS_PRE, "{\"schema\":\"sai-eventchange\", " "\"event_hash\":\"%s\", \"state\":%d}", - arg->uid, arg->state); + event_uuid, state); memset(&info, 0, sizeof(info)); info.private_source_idx = SAI_WEBSRV_PB__GENERATED; @@ -197,21 +182,11 @@ _sais_eventchange(struct lws_ss_handle *h, void *_arg) info.len = (size_t)n; info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, &info) < 0) { + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(hsrv, &info) < 0) { lwsl_warn("%s: buflist append failed\n", __func__); + return; } - - if (lws_ss_request_tx(h)) - lwsl_ss_warn(h, "req fail"); -} - -void -sais_eventchange(struct lws_ss_handle *hsrv, const char *event_uuid, int state) -{ - struct sais_arg arg = { event_uuid, state }; - - lws_ss_server_foreach_client(hsrv, _sais_eventchange, (void *)&arg); } sai_db_result_t @@ -362,7 +337,7 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) return SAI_DB_RESULT_ERROR; lws_sql_purify(esc, platform, sizeof(esc)); - lws_snprintf(filt, sizeof(filt), " and platform='%s' and state=4", esc); + lws_snprintf(filt, sizeof(filt), " and platform='%s'", esc); if (lws_struct_sq3_deserialize(pdb, filt, NULL, lsm_schema_sq3_map_task, diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c deleted file mode 100644 index d5017f9..0000000 --- a/src/server/s-websrv.c +++ /dev/null @@ -1,669 +0,0 @@ -/* - * Sai server - * - * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - * - * - * This is a ws server over a unix domain socket made available by sai-server - * and connected to by sai-web instances running on the same box. - * - * The server notifies the sai-web instances of event and task changes (just - * that a particular event or task changed) and builder list updates (the - * whole current builder list JSON each time). - * - * Sai-web instances can send requests to restart or delete tasks and whole - * events made by authenticated clients. - * - * Since this is on a local UDS protected by user:group, there's no tls or auth - * on this link itself. - */ - -#include <libwebsockets.h> -#include <string.h> -#include <signal.h> -#include <time.h> -#include <assert.h> - -#include "s-private.h" - -typedef struct sai_sul_retry_ctx { - lws_sorted_usec_list_t sul; - struct vhd *vhd; - char uuid[SAI_TASKID_LEN + 1]; - char platform[96]; - int retries; - uint8_t op; /* SAIS_WS_WEBSRV_RX_... */ -} sai_sul_retry_ctx_t; - - -static lws_struct_map_t lsm_browser_taskreset[] = { - LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "uuid"), -}; - -static lws_struct_map_t lsm_browser_platreset[] = { - LSM_CARRAY (sai_browse_rx_platreset_t, event_uuid, "event_uuid"), - LSM_CARRAY (sai_browse_rx_platreset_t, platform, "platform"), -}; - -static const lws_struct_map_t lsm_viewercount_members[] = { - LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"), -}; - -static const lws_struct_map_t lsm_schema_json_map[] = { - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, - /* shares struct */ "com.warmcat.sai.taskreset"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, - /* shares struct */ "com.warmcat.sai.taskrebuildlaststep"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, - /* shares struct */ "com.warmcat.sai.eventreset"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, - /* shares struct */ "com.warmcat.sai.eventdelete"), - LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, - "com.warmcat.sai.taskcan"), - LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewercount_members, - "com.warmcat.sai.viewercount"), - LSM_SCHEMA (sai_rebuild_t, NULL, lsm_rebuild, - "com.warmcat.sai.rebuild"), - LSM_SCHEMA (sai_browse_rx_platreset_t, NULL, lsm_browser_platreset, - "com.warmcat.sai.platreset"), - LSM_SCHEMA (sai_stay_t, NULL, lsm_stay, - "com.warmcat.sai.stay"), -}; - -enum { - SAIS_WS_WEBSRV_RX_TASKRESET, - SAIS_WS_WEBSRV_RX_TASKREBUILDLASTSTEP, - SAIS_WS_WEBSRV_RX_EVENTRESET, - SAIS_WS_WEBSRV_RX_EVENTDELETE, - SAIS_WS_WEBSRV_RX_TASKCANCEL, - SAIS_WS_WEBSRV_RX_VIEWERCOUNT, - SAIS_WS_WEBSRV_RX_REBUILD, - SAIS_WS_WEBSRV_RX_PLATRESET, - SAIS_WS_WEBSRV_RX_STAY, -}; - -void -sais_mark_all_builders_offline(struct vhd *vhd) -{ - char *err = NULL; - - lwsl_notice("%s: marking all builders offline initially\n", __func__); - - sqlite3_exec(vhd->server.pdb, "UPDATE builders SET online = 0;", - NULL, NULL, &err); - if (err) { - lwsl_err("%s: sqlite error: %s\n", __func__, err); - sqlite3_free(err); - } -} - -int -sais_validate_id(const char *id, int reqlen) -{ - const char *idin = id; - int n = reqlen; - - while (*id && n--) { - if (!((*id >= '0' && *id <= '9') || - (*id >= 'a' && *id <= 'z') || - (*id >= 'A' && *id <= 'Z'))) - goto reject; - id++; - } - - if (!n && !*id) - return 0; -reject: - - lwsl_notice("%s: Invalid ID (%d) '%s'\n", __func__, reqlen, idin); - - return 1; -} - -static int -sais_validate_builder_name(const char *id) -{ - const char *idin = id; - - while (*id) { - if (!((*id >= '0' && *id <= '9') || - (*id >= 'a' && *id <= 'z') || - (*id >= 'A' && *id <= 'Z') || - *id == '.' || *id == '/' || *id == '-' || *id == '_')) - goto reject; - id++; - } - - return 0; -reject: - - lwsl_notice("%s: Invalid builder name '%s'\n", __func__, idin); - - return 1; -} - - -int -sais_list_builders(struct vhd *vhd) -{ - char json_builders[LWS_PRE + 1024], *start = json_builders + LWS_PRE, - *p = start, *end = p + sizeof(json_builders) - LWS_PRE, - subsequent = 0; - unsigned int ss_flags = LWSSS_FLAG_SOM; - lws_dll2_owner_t db_builders_owner; - sai_plat_t *builder_from_db; - lws_struct_serialize_t *js; - struct lwsac *ac = NULL; - lws_wsmsg_info_t info; - size_t w; - - memset(&db_builders_owner, 0, sizeof(db_builders_owner)); - - if (lws_struct_sq3_deserialize(vhd->server.pdb, NULL, "name ", - lsm_schema_sq3_map_plat, - &db_builders_owner, &ac, 0, 100)) { - lwsl_err("%s: Failed to query builders from DB\n", __func__); - return 1; - } - - // lwsl_warn("%s: count deserialized %d\n", __func__, (int)db_builders_owner.count); - - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), - "{\"schema\":\"com.warmcat.sai.builders\",\"builders\":["); - - lws_start_foreach_dll(struct lws_dll2 *, walk, db_builders_owner.head) { - lws_struct_json_serialize_result_t r; - sai_plat_t *live_builder; - - builder_from_db = lws_container_of(walk, sai_plat_t, sai_plat_list); - - /* - * Find this builder in the live list by name. This is safe because - * builder_from_db->name is a valid string within the scope of this function. - */ - live_builder = sais_builder_from_uuid(vhd, builder_from_db->name, __FILE__, __LINE__); - - if (live_builder) { - // lwsl_notice("%s: live_builder %s found, stay_on: %d, copying to db_builder (stay_on: %d)\n", - // __func__, live_builder->name, live_builder->stay_on, builder_from_db->stay_on); - builder_from_db->online = 1; - lws_strncpy(builder_from_db->peer_ip, live_builder->peer_ip, - sizeof(builder_from_db->peer_ip)); - builder_from_db->stay_on = live_builder->stay_on; - } else - builder_from_db->online = 0; - - builder_from_db->powering_up = 0; - builder_from_db->powering_down = 0; - - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.power_state_owner.head) { - sai_power_state_t *ps = lws_container_of(p, sai_power_state_t, list); - size_t host_len = strlen(ps->host); - - if (!strncmp(builder_from_db->name, ps->host, host_len) && - builder_from_db->name[host_len] == '.') { - builder_from_db->powering_up = ps->powering_up; - builder_from_db->powering_down = ps->powering_down; - break; - } - } lws_end_foreach_dll(p); - - js = lws_struct_json_serialize_create(lsm_schema_map_plat_simple, - LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), - 0, builder_from_db); - if (!js) - goto bail; - - if (subsequent) - *p++ = ','; - subsequent = 1; - - do { - r = lws_struct_json_serialize(js, (uint8_t *)p, - lws_ptr_diff_size_t(end, p) - 2, &w); - p += w; - - switch (r) { - case LSJS_RESULT_FINISH: - /* fallthru */ - case LSJS_RESULT_CONTINUE: - memset(&info, 0, sizeof(info)); - - info.private_source_idx = SAI_WEBSRV_PB__GENERATED; - info.buf = (uint8_t *)start; - info.len = lws_ptr_diff_size_t(p, start); - info.ss_flags = ss_flags; - - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) - lwsl_warn("%s: unable to broadcast to web\n", __func__); - - p = start; - ss_flags &= ~((unsigned int)LWSSS_FLAG_SOM); - break; - - case LSJS_RESULT_ERROR: - lws_struct_json_serialize_destroy(&js); - goto bail; - } - - } while (r == LSJS_RESULT_CONTINUE); - - lws_struct_json_serialize_destroy(&js); - - } lws_end_foreach_dll(walk); - - ss_flags |= LWSSS_FLAG_EOM; - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - memset(&info, 0, sizeof(info)); - - info.private_source_idx = SAI_WEBSRV_PB__GENERATED; - info.buf = (uint8_t *)start; - info.len = lws_ptr_diff_size_t(p, start); - info.ss_flags = ss_flags; - - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) - lwsl_warn("%s: unable to broadcast to web\n", __func__); - - // lwsl_notice("%s: Broadcasting builder list: %s\n", __func__, start); - lwsac_free(&ac); - return 0; - -bail: - lwsac_free(&ac); - return 1; -} - - - -static void -sum_viewers_cb(struct lws_ss_handle *h, void *arg) -{ - websrvss_srv_t *m_client = (websrvss_srv_t *)lws_ss_to_user_object(h); - *(unsigned int *)arg += m_client->viewers; -} - - - - -static lws_ss_state_return_t -websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) -{ - websrvss_srv_t *m = (websrvss_srv_t *)userobj; - sai_browse_rx_evinfo_t *ei; - lws_struct_args_t a; - int n; - - // lwsl_user("%s: len %d, flags: %d\n", __func__, (int)len, flags); - // lwsl_hexdump_info(buf, len); - - memset(&a, 0, sizeof(a)); - a.map_st[0] = lsm_schema_json_map; - a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_json_map); - a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_json_map); - a.ac_block_size = 128; - - lws_struct_json_init_parse(&m->ctx, NULL, &a); - n = lejp_parse(&m->ctx, (uint8_t *)buf, (int)len); - if (n < 0 || !a.dest) { - lwsl_hexdump_notice(buf, len); - lwsl_notice("%s: notification JSON decode failed '%s'\n", - __func__, lejp_error_to_string(n)); - return LWSSSSRET_DISCONNECT_ME; - } - - // lwsl_notice("%s: schema idx %d\n", __func__, a.top_schema_index); - - switch (a.top_schema_index) { - case SAIS_WS_WEBSRV_RX_TASKRESET: - { - 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_TASKRESET: %s: received", ei->event_hash); - if (sais_task_clear_build_and_logs(m->vhd, ei->event_hash, 0)) - lwsl_ss_err(m->ss, "taskreset failed"); - 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)) - goto soft_error; - - lwsl_ss_warn(m->ss, "SAIS_WS_WEBSRV_RX_TASKREBUILDLASTSTEP: %s: received", ei->event_hash); - if (sais_task_rebuild_last_step(m->vhd, ei->event_hash)) - lwsl_ss_err(m->ss, "taskrebuildlaststep failed"); - break; - } - - case SAIS_WS_WEBSRV_RX_EVENTRESET: - { - sai_db_result_t r; - - ei = (sai_browse_rx_evinfo_t *)a.dest; - - if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) - goto soft_error; - - r = sais_event_reset(m->vhd, ei->event_hash); - if (r) - lwsl_ss_err(m->ss, "eventreset failed"); - - lwsac_free(&a.ac); - break; - } - - case SAIS_WS_WEBSRV_RX_PLATRESET: { - sai_browse_rx_platreset_t *pr = (sai_browse_rx_platreset_t *)a.dest; - sai_db_result_t r; - - if (sais_validate_id(pr->event_uuid, SAI_EVENTID_LEN)) - goto soft_error; - - r = sais_plat_reset(m->vhd, pr->event_uuid, pr->platform); - if (r) - lwsl_ss_err(m->ss, "platreset failed"); - lwsac_free(&a.ac); - break; - } - - case SAIS_WS_WEBSRV_RX_EVENTDELETE: - { - sai_db_result_t r; - - ei = (sai_browse_rx_evinfo_t *)a.dest; - if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) { - lwsl_err("%s: SAIS_WS_WEBSRV_RX_EVENTDELETE: unable to validate id %s\n", __func__, ei->event_hash); - goto soft_error; - } - - lwsl_notice("%s: eventdelete %s\n", __func__, ei->event_hash); - - r = sais_event_delete(m->vhd, ei->event_hash); - if (r) - lwsl_ss_err(m->ss, "event delete failed"); - lwsac_free(&a.ac); - break; - } - - case SAIS_WS_WEBSRV_RX_TASKCANCEL: - - - ei = (sai_browse_rx_evinfo_t *)a.dest; - if (sais_validate_id(ei->event_hash, SAI_TASKID_LEN)) - goto soft_error; - - sais_task_cancel(m->vhd, ei->event_hash); - - break; - - case SAIS_WS_WEBSRV_RX_VIEWERCOUNT: - { - sai_viewer_state_t *vs = (sai_viewer_state_t *)a.dest; - unsigned int total_viewers = 0; - char old_viewers_present = !!m->vhd->viewers_are_present; - - /* Store viewer count for this specific sai-web client */ - m->viewers = vs->viewers; - - /* Recalculate total from all connected sai-web clients */ - lws_ss_server_foreach_client(m->vhd->h_ss_websrv, - sum_viewers_cb, &total_viewers); - - m->vhd->browser_viewer_count = total_viewers; - - m->vhd->viewers_are_present = !!total_viewers; - - /* - * Only broadcast to builders if the state has changed - * from 0 viewers to >0, or from >0 viewers to 0. - */ - if (old_viewers_present != m->vhd->viewers_are_present) { - lwsl_notice("%s: Viewer presence changed to %d. Broadcasting to builders.\n", - __func__, m->vhd->viewers_are_present); - lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) { - struct pss *pss_builder = lws_container_of(p, struct pss, same); - sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); - - if (vsend) { - vsend->viewers = m->vhd->viewers_are_present; - lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); - lws_callback_on_writable(pss_builder->wsi); - } - } lws_end_foreach_dll(p); - } - break; - } - case SAIS_WS_WEBSRV_RX_REBUILD: - { - sai_rebuild_t *reb = (sai_rebuild_t *)a.dest; - sai_plat_t *cb; - - if (sais_validate_builder_name(reb->builder_name)) - goto soft_error; - - cb = sais_builder_from_uuid(m->vhd, reb->builder_name, - __FILE__, __LINE__); - if (!cb) { - lwsl_info("%s: unknown builder %s for rebuild\n", - __func__, reb->builder_name); - lwsac_free(&a.ac); - break; - } - - /* cb->wsi is the builder connection */ - lws_start_foreach_dll(struct lws_dll2 *, p, - m->vhd->builders.head) { - struct pss *pss = lws_container_of(p, struct pss, same); - - if (pss->wsi == cb->wsi) { - sai_rebuild_t *r = malloc(sizeof(*r)); - if (!r) - break; - *r = *reb; - lws_dll2_add_tail(&r->list, - &pss->rebuild_owner); - lws_callback_on_writable(pss->wsi); - break; - } - } lws_end_foreach_dll(p); - } - break; - - case SAIS_WS_WEBSRV_RX_STAY: - { - sai_stay_t *stay = (sai_stay_t *)a.dest; - - lws_start_foreach_dll(struct lws_dll2 *, p, - m->vhd->sai_powers.head) { - struct pss *pss_power = lws_container_of(p, struct pss, same); - sai_stay_t *s; - - s = malloc(sizeof(*s)); - if (s) { - *s = *stay; - lws_dll2_add_tail(&s->list, &pss_power->stay_owner); - lws_callback_on_writable(pss_power->wsi); - } - } lws_end_foreach_dll(p); - - lwsac_free(&a.ac); - break; - } - } - - return 0; - -soft_error: - lwsl_warn("%s: soft error\n", __func__); - - return 0; -} - -static lws_ss_state_return_t -websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, - size_t *len, int *flags) -{ - websrvss_srv_t *m = (websrvss_srv_t *)userobj; - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->bl_srv_to_web), depi, fl; - char som, som1, eom, final = 1; - size_t fsl, used; - - if (!m->bl_srv_to_web) - return LWSSSSRET_TX_DONT_SEND; - - depi = *pi; - - /* - * We can only issue *len at a time. - * - * Notice we are getting the stored flags from the START of the fragment each time. - * that means we can still see the right flags stored with the fragment, even if we - * have partially used the buflist frag and are partway through it. - * - * Ergo, only something to skip if we are at som=1. And also notice that although - * *pi will be right, after the lws_buflist..._use() api, what it points to has been - * destroyed. So we also dereference *pi into depi for use below. - */ - - fsl = lws_buflist_next_segment_len(&m->bl_srv_to_web, NULL); - - lws_buflist_fragment_use(&m->bl_srv_to_web, NULL, 0, &som, &eom); - if (som) { - fsl -= sizeof(int); - lws_buflist_fragment_use(&m->bl_srv_to_web, buf, sizeof(int), &som1, &eom); - } - if (!(depi & LWSSS_FLAG_SOM)) - som = 0; - - used = (size_t)lws_buflist_fragment_use(&m->bl_srv_to_web, (uint8_t *)buf, *len, &som1, &eom); - if (!used) - return LWSSSSRET_TX_DONT_SEND; - - if (used < fsl || !(depi & LWSSS_FLAG_EOM)) /* we saved SS flags at the start of the buf */ - final = 0; - - *len = used; - fl = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - - // lwsl_ss_notice(m->ss, "Sending %d srv->web: ssflags %d", (int)*len, fl); - - if ((fl & LWSSS_FLAG_SOM) && (((*flags) & 3) == 2)) { - lwsl_ss_err(m->ss, "TX: Illegal LWSSS_FLAG_SOM after previous frame without LWSSS_FLAG_EOM"); - assert(0); - } - if (!(fl & LWSSS_FLAG_SOM) && ((*flags) & 3) == 3) { - lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM after previous frame with LWSSS_FLAG_EOM"); - assert(0); - } - if (!(fl & LWSSS_FLAG_SOM) && !((*flags) & 2)) { - lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM on first frame"); - assert(0); - } - - *flags = fl; - - - // lwsl_hexdump_notice(buf, *len); - - if (m->bl_srv_to_web) - return lws_ss_request_tx(m->ss); - - return 0; -} - - -static lws_ss_state_return_t -websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, - lws_ss_tx_ordinal_t ack) -{ - websrvss_srv_t *m = (websrvss_srv_t *)userobj; - - // lwsl_user("%s: %p %s, ord 0x%x\n", __func__, m->ss, - // lws_ss_state_name((int)state), (unsigned int)ack); - - switch (state) { - case LWSSSCS_DISCONNECTED: { - unsigned int total_viewers = 0; - - lws_buflist_destroy_all_segments(&m->bl_srv_to_web); - lws_wsmsg_destroy(m->private_heads, LWS_ARRAY_SIZE(m->private_heads)); - - m->viewers = 0; - - /* This sai-web client disconnected, recalculate total viewers */ - lws_ss_server_foreach_client(m->vhd->h_ss_websrv, - sum_viewers_cb, &total_viewers); - - m->vhd->browser_viewer_count = total_viewers; - char new_viewers_present = !!total_viewers; - - if (m->vhd->viewers_are_present != new_viewers_present) { - m->vhd->viewers_are_present = !!new_viewers_present; - lwsl_notice("%s: A sai-web client disconnected, viewer presence changed to %d. Broadcasting.\n", - __func__, new_viewers_present); - - /* Broadcast new presence state to builders */ - lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) { - struct pss *pss_builder = lws_container_of(p, struct pss, same); - sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); - - if (vsend) { - vsend->viewers = (unsigned int)new_viewers_present; - lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); - lws_callback_on_writable(pss_builder->wsi); - } - } lws_end_foreach_dll(p); - } - - break; - } - case LWSSSCS_CREATING: - m->viewers = 0; - return lws_ss_request_tx(m->ss); - - case LWSSSCS_CONNECTED: - sais_list_builders(m->vhd); - break; - case LWSSSCS_ALL_RETRIES_FAILED: - break; - - case LWSSSCS_SERVER_TXN: - break; - - case LWSSSCS_SERVER_UPGRADE: - break; - - default: - break; - } - - return 0; -} - -const lws_ss_info_t ssi_server = { - .handle_offset = offsetof(websrvss_srv_t, ss), - .opaque_user_data_offset = offsetof(websrvss_srv_t, vhd), - .streamtype = "websrv", - .rx = websrvss_ws_rx, - .tx = websrvss_ws_tx, - .state = websrvss_srv_state, - .user_alloc = sizeof(websrvss_srv_t), -}; diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index bfc8581..ae51ea5 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -20,6 +20,10 @@ * * These are ws rx and tx handlers related to builder ws connections, at the * sai-server + * + * b1 --\ sai- sai- /-- browser + * b2 ----- server ---- web ------ browser + * b3 --/ * \-- browser */ #include <libwebsockets.h> @@ -29,7 +33,6 @@ #include <time.h> #include "s-private.h" -#include "s-metrics-db.h" const lws_struct_map_t lsm_schema_map_ta[] = { LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), @@ -178,8 +181,11 @@ sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) static void sais_log_to_db(struct vhd *vhd, sai_log_t *log) { + char event_uuid[33], q[256], esc_uuid[129]; sais_logcache_pertask_t *lcpt = NULL; + sqlite3 *pdb = NULL; sai_log_t *hlog; + int step; /* * find the pertask if one exists @@ -223,50 +229,44 @@ sais_log_to_db(struct vhd *vhd, sai_log_t *log) if (!vhd->sul_logcache.list.owner) /* if not already scheduled, schedule it for 250ms */ - lws_sul_schedule(vhd->context, 0, &vhd->sul_logcache, sais_dump_logs_to_db, 250 * LWS_US_PER_MS); - if (log->channel == 3 && log->log) { /* control channel */ - int step; - if (!memcmp(log->log, " Step ", 5)) { - char event_uuid[33]; - sqlite3 *pdb = NULL; - char q[256], esc_uuid[129]; + if (log->channel != 3 /* control channel */ || !log->log || + log->len < 5 || memcmp(log->log, " Step ", 5)) + return; - step = atoi(&log->log[5]); + step = atoi(&log->log[5]); - sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); - if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); - lws_sql_purify(esc_uuid, log->task_uuid, sizeof(esc_uuid)); + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + return; - lws_snprintf(q, sizeof(q), - "UPDATE tasks SET build_step=%d WHERE uuid='%s'", - step, esc_uuid); - if (sai_sqlite3_statement(pdb, q, "update build_step")) - lwsl_err("%s: failed to update build_step\n", __func__); + lws_sql_purify(esc_uuid, log->task_uuid, sizeof(esc_uuid)); - sais_event_db_close(vhd, &pdb); + lws_snprintf(q, sizeof(q), + "UPDATE tasks SET build_step=%d WHERE uuid='%s'", + step, esc_uuid); - sais_taskchange(vhd->h_ss_websrv, log->task_uuid, SAIES_BEING_BUILT); - } - } - } + if (sai_sqlite3_statement(pdb, q, "update build_step")) + lwsl_err("%s: failed to update build_step\n", __func__); + + sais_event_db_close(vhd, &pdb); } sai_plat_t * -sais_builder_from_uuid(struct vhd *vhd, const char *hostname, const char *_file, int _line) +sais_builder_from_uuid(struct vhd *vhd, const char *hostname) { lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - sai_plat_t *cb = lws_container_of(p, sai_plat_t, + sai_plat_t *sp = lws_container_of(p, sai_plat_t, sai_plat_list); - if (!strcmp(hostname, cb->name)) { - lwsl_info("%s: %s:%d: found live builder %s\n", __func__, _file, _line, hostname); - cb->online = 1; - return cb; + if (!strcmp(hostname, sp->name)) { + sp->online = 1; + + return sp; } } lws_end_foreach_dll(p); @@ -279,13 +279,13 @@ sais_builder_from_host(struct vhd *vhd, const char *host) { lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - sai_plat_t *cb = lws_container_of(p, sai_plat_t, + sai_plat_t *sp = lws_container_of(p, sai_plat_t, sai_plat_list); size_t host_len = strlen(host); - if (!strncmp(cb->name, host, host_len) && - cb->name[host_len] == '.') - return cb; + if (!strncmp(sp->name, host, host_len) && + sp->name[host_len] == '.') + return sp; } lws_end_foreach_dll(p); @@ -296,22 +296,35 @@ void sais_set_builder_power_state(struct vhd *vhd, const char *name, int up, int down) { sai_power_state_t *ps = NULL; - sai_plat_t *live_builder; + sai_plat_t *live_builder = sais_builder_from_host(vhd, name); - if (up) { - live_builder = sais_builder_from_host(vhd, name); - if (live_builder) - return; - } + if (live_builder && up) + up = 0; - if (down) { - live_builder = sais_builder_from_host(vhd, name); - if (!live_builder) - return; - } + if (!live_builder && down) + down = 0; + + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + vhd->server.power_state_owner.head) { + ps = lws_container_of(p, sai_power_state_t, list); + + if (!strcmp(ps->host, name)) { + if (live_builder && ps->powering_up) + ps->powering_up = 0; + + if (!live_builder && ps->powering_down) + ps->powering_down = 0; + + if (!ps->powering_up && !ps->powering_down) { + lws_dll2_remove(&ps->list); + free(ps); + } + } + } lws_end_foreach_dll_safe(p, p1); lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.power_state_owner.head) { ps = lws_container_of(p, sai_power_state_t, list); + if (!strcmp(ps->host, name)) break; ps = NULL; @@ -335,7 +348,7 @@ sais_set_builder_power_state(struct vhd *vhd, const char *name, int up, int down } } - sais_list_builders(vhd); + sais_list_builders(vhd); } /* @@ -344,9 +357,9 @@ sais_set_builder_power_state(struct vhd *vhd, const char *name, int up, int down void sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) { - sai_plat_t *cb; struct lwsac *ac = NULL; lws_dll2_owner_t o; + sai_plat_t *sp; int n; /* @@ -356,13 +369,13 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) */ lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, vhd->server.builder_owner.head) { - cb = lws_container_of(p, sai_plat_t, sai_plat_list); + sp = lws_container_of(p, sai_plat_t, sai_plat_list); - if (cb->wsi == wsi) { + if (sp->wsi == wsi) { char q[256]; lwsl_notice("%s: Builder '%s' disconnected\n", __func__, - cb->name); + sp->name); /* * Check all active events for tasks that were running @@ -388,12 +401,12 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) SAIES_BEING_BUILT); if (sqlite3_prepare_v2(pdb, q, -1, &sm, NULL) == SQLITE_OK) { - sqlite3_bind_text(sm, 1, cb->name, -1, SQLITE_TRANSIENT); + sqlite3_bind_text(sm, 1, sp->name, -1, SQLITE_TRANSIENT); while (sqlite3_step(sm) == SQLITE_ROW) { const unsigned char *task_uuid = sqlite3_column_text(sm, 0); if (task_uuid) { lwsl_notice("%s: resetting task %s from disconnected builder %s\n", - __func__, (const char *)task_uuid, cb->name); + __func__, (const char *)task_uuid, sp->name); sais_task_clear_build_and_logs(vhd, (const char *)task_uuid, 0); } } @@ -409,7 +422,7 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) /* drop any inflight task information for this builder */ lws_start_foreach_dll_safe(struct lws_dll2 *, pif, pif1, - cb->inflight_owner.head) { + sp->inflight_owner.head) { sai_uuid_list_t *ul = lws_container_of(pif, sai_uuid_list_t, list); sais_inflight_entry_destroy(ul); @@ -417,10 +430,10 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) } lws_end_foreach_dll_safe(pif, pif1); - const char *dot = strchr(cb->name, '.'); + const char *dot = strchr(sp->name, '.'); if (dot) { char host[128]; - lws_strnncpy(host, cb->name, dot - cb->name, sizeof(host)); + lws_strnncpy(host, sp->name, dot - sp->name, sizeof(host)); lws_start_foreach_dll_safe(struct lws_dll2 *, p2, p3, vhd->server.power_state_owner.head) { sai_power_state_t *ps = lws_container_of(p2, sai_power_state_t, list); if (!strcmp(ps->host, host)) { @@ -431,9 +444,9 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) } lws_end_foreach_dll_safe(p2, p3); } - lws_dll2_remove(&cb->sai_plat_list); - lws_sul_cancel(&cb->sul_find_jobs); - free(cb); + lws_dll2_remove(&sp->sai_plat_list); + lws_sul_cancel(&sp->sul_find_jobs); + free(sp); // assert(0); } @@ -462,19 +475,22 @@ int sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl, unsigned int ss_flags) { char event_uuid[33], s[128], esc[96], do_remove_uuid; - const sai_build_metric_t *metric; sai_resource_requisition_t *rr; sai_resource_wellknown_t *wk; sai_plat_owner_t *bp_owner; + lws_struct_serialize_t *js; + sai_build_metric_t *metric; struct lwsac *ac = NULL; - sai_plat_t *build, *cb; + sai_plat_t *build, *sp; lws_wsmsg_info_t info; sai_rejection_t *rej; sai_resource_t *res; sai_uuid_list_t *ul; lws_dll2_owner_t o; sai_artifact_t *ap; + uint8_t xbuf[2048]; sai_task_t *task; + size_t used = 0; sai_log_t *log; uint64_t rid; int n, m; @@ -542,7 +558,13 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b * with several builders connected and spamming * fragmented load reports, when we forward them * the adjacent fragments will be randomly - * ordered. Even though each builder is sending + * ordered (* shows where this code is) + * + * b1 --\ sai- sai- /-- browser + * b2 ----- server ---- web ------ browser + * b3 --/ * \-- browser + * + * Even though each builder is sending * them correctly ordered, when all combined * together on the srv -> web link, the fragments * will be disorderd. Eg, b1 first frag, b2 @@ -611,10 +633,8 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b /* * Step 2: Update the long-lived, malloc'd in-memory list. */ - //cb = sais_builder_from_uuid(vhd, build->name); - //if (cb) - // sais_builder_disconnected(vhd, cb->wsi); - live_sp = sais_builder_from_uuid(vhd, build->name, __FILE__, __LINE__); + + live_sp = sais_builder_from_uuid(vhd, build->name); if (live_sp) { /* Already exists (reconnect), just update dynamic info */ lwsl_err("%s: found live builder for %s\n", __func__, build->name); @@ -628,12 +648,8 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b sizeof(live_sp->lws_hash)); live_sp->windows = build->windows; live_sp->online = 1; - live_sp->avail_slots = -1; /* ie, unknown */ live_sp->avail_mem_kib = (unsigned int)-1; live_sp->avail_sto_kib = (unsigned int)-1; - live_sp->s_avail_slots = live_sp->avail_slots; - live_sp->s_inflight_count = (int)live_sp->inflight_owner.count; - live_sp->s_last_rej_task_uuid[0] = '\0'; } else { /* New builder, create a deep-copied, malloc'd object */ size_t nlen = strlen(build->name) + 1; @@ -655,21 +671,20 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b lws_strncpy(live_sp->lws_hash, build->lws_hash, sizeof(live_sp->lws_hash)); live_sp->windows = build->windows; - live_sp->avail_slots = 1; /* default */ live_sp->avail_mem_kib = (unsigned int)-1; live_sp->avail_sto_kib = (unsigned int)-1; - live_sp->s_avail_slots = live_sp->avail_slots; live_sp->wsi = pss->wsi; live_sp->cx = lws_get_context(pss->wsi); live_sp->vhd = vhd; live_sp->online = 1; lws_strncpy(live_sp->peer_ip, pss->peer_ip, sizeof(live_sp->peer_ip)); + lws_dll2_add_tail(&live_sp->sai_plat_list, &vhd->server.builder_owner); } } lws_sul_schedule(live_sp->cx, 0, &live_sp->sul_find_jobs, - sais_plat_find_jobs_cb, 1 * LWS_US_PER_SEC); + sais_plat_find_jobs_cb, 500 * LWS_US_PER_MS); const char *dot = strchr(build->name, '.'); if (dot) { @@ -688,19 +703,19 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b * builder that just connected. */ lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - cb = lws_container_of(p, sai_plat_t, sai_plat_list); - if (cb->wsi == pss->wsi) { + sp = lws_container_of(p, sai_plat_t, sai_plat_list); + if (sp->wsi == pss->wsi) { /* This platform belongs to the connection that sent the message */ - if (sais_allocate_task(vhd, pss, cb, cb->platform) < 0) + if (sais_allocate_task(vhd, pss, sp, sp->platform) < 0) goto bail; } } lws_end_foreach_dll(p); #if 0 lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - cb = lws_container_of(p, sai_plat_t, sai_plat_list); - if (cb->wsi == pss->wsi) { + sp = lws_container_of(p, sai_plat_t, sai_plat_list); + if (sp->wsi == pss->wsi) { /* This platform belongs to the connection that sent the message */ - if (sais_allocate_task(vhd, pss, cb, cb->platform) < 0) + if (sais_allocate_task(vhd, pss, sp, sp->platform) < 0) goto bail; } } lws_end_foreach_dll(p); @@ -727,93 +742,6 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b log = (sai_log_t *)pss->a.dest; sais_log_to_db(vhd, log); - if (pss->mark_started) { - pss->mark_started = 0; - pss->first_log_timestamp = log->timestamp; - // if (sais_set_task_state(vhd, NULL, NULL, log->task_uuid, - // SAIES_BEING_BUILT, 0, 0)) - // goto bail; - // sais_create_and_offer_task_step(vhd, log->task_uuid, 11); - } - - - if (log->finished) { - sai_plat_t *cb; - // sai_uuid_list_t *u; - char builder_name[128], esc_uuid[129], q[128], event_uuid[33]; - sqlite3 *pdb = NULL; - - /* - * This step is finished, find the builder. - * - * We don't move on its state until we receive the - * SAI_TASK_REASON_ from the "_REJ" message from the - * builder. - */ - - sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); - if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { - builder_name[0] = '\0'; - lws_sql_purify(esc_uuid, log->task_uuid, sizeof(esc_uuid)); - lws_snprintf(q, sizeof(q), - "select builder_name from tasks where uuid='%s'", - esc_uuid); - if (sqlite3_exec(pdb, q, sql3_get_string_cb, builder_name, - NULL) == SQLITE_OK && builder_name[0]) { - cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); - if (cb) { - // sai_uuid_list_t *sul; - - lwsl_notice("%s: builder %s reports step done, slots %d, mem %d, sto %d\n", - __func__, cb->name, log->avail_slots, log->avail_mem_kib, log->avail_sto_kib); - - cb->avail_slots = log->avail_slots; - cb->avail_mem_kib = log->avail_mem_kib; - cb->avail_sto_kib = log->avail_sto_kib; - cb->last_rej_task_uuid[0] = '\0'; - - cb->s_avail_slots = cb->avail_slots; - cb->s_inflight_count = (int)cb->inflight_owner.count; - lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, - sizeof(cb->s_last_rej_task_uuid)); - sais_list_builders(vhd); - } - } - sais_event_db_close(vhd, &pdb); - } - - /* - * We have reached the end of the logs for this task step - */ - - sais_dump_logs_to_db(&vhd->sul_logcache); - - lwsl_notice("%s: \\\\\\\\\\\\\\\\\\ log->finished says 0x%x, dur %lluus\n", - __func__, log->finished, (unsigned long long)( - log->timestamp - pss->first_log_timestamp)); - if (log->finished & SAISPRF_EXIT) { - if ((log->finished & 0xff) == 0) { - n = SAIES_STEP_SUCCESS; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_SUCCESS: %s\n", __func__, log->task_uuid); - } else { - n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL: %s\n", __func__, log->task_uuid); - } - } else - if (log->finished & 0x2000) { - n = SAIES_CANCELLED; - lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED: %s\n", __func__, log->task_uuid); - - } else { - n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_FAIL: %s\n", __func__, log->task_uuid); - } - - if (sais_set_task_state(vhd, NULL, NULL, log->task_uuid, n, 0, - log->timestamp - pss->first_log_timestamp)) - goto bail; - } - lwsac_free(&pss->a.ac); break; @@ -830,21 +758,17 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b break; rej->host_platform[sizeof(rej->host_platform) - 1] = '\0'; - cb = sais_builder_from_uuid(vhd, rej->host_platform, __FILE__, __LINE__); - if (!cb) { + sp = sais_builder_from_uuid(vhd, rej->host_platform); + if (!sp) { lwsl_info("%s: unknown builder %s rejecting\n", __func__, rej->host_platform); lwsac_free(&pss->a.ac); break; } - cb->avail_slots = rej->avail_slots; - cb->avail_mem_kib = rej->avail_mem_kib; - cb->avail_sto_kib = rej->avail_sto_kib; - lwsl_notice("%s: builder %s reports task status update, reason: %d, %s, slots %d, mem %d, sto %d\n", - __func__, cb->name, rej->reason, rej->task_uuid, - cb->avail_slots, cb->avail_mem_kib, cb->avail_sto_kib); + __func__, sp->name, rej->reason, rej->task_uuid, + sp->avail_slots, sp->avail_mem_kib, sp->avail_sto_kib); do_remove_uuid = 0; @@ -865,38 +789,86 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b lws_snprintf(q, sizeof(q), "select build_step from tasks where uuid='%s'", esc_uuid); + if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &build_step, NULL) != SQLITE_OK) build_step = -1; + + /* + * Bump the build step on the accepted task + */ + + build_step++; + lws_snprintf(q, sizeof(q), "update tasks set build_step=%d where state != 4 and uuid='%s'", + build_step, esc_uuid); + sqlite3_exec(pdb, q, NULL, NULL, NULL); + + if (build_step == 1) { + pss->first_log_timestamp = (uint64_t)lws_now_secs(); + lws_snprintf(q, sizeof(q), + "update tasks set started=%llu where uuid='%s'", + (unsigned long long)pss->first_log_timestamp, esc_uuid); + + if (sqlite3_exec(pdb, q, NULL, NULL, NULL) != SQLITE_OK) + lwsl_notice("%s: unable to set started\n", __func__); + } + + lwsl_notice("%s: exiting, setting build_step %d\n", __func__, build_step); + sais_event_db_close(vhd, &pdb); - } - if (build_step == 0) - pss->first_log_timestamp = (uint64_t)lws_now_usecs(); + if (sais_set_task_state(vhd, + rej->task_uuid, + SAIES_BEING_BUILT, + !build_step ? pss->first_log_timestamp : 0, 0)) + break; + } } - - if (sais_set_task_state(vhd, NULL, NULL, rej->task_uuid, - SAIES_BEING_BUILT, 0, 0)) - break; /* leave the uuid listed as inflight until step completed */ break; + case SAI_TASK_REASON_DUPE: lwsl_notice("%s: SAI_TASK_REASON_DUPE: %s\n", __func__, rej->task_uuid); break; + case SAI_TASK_REASON_BUSY: lwsl_notice("%s: SAI_TASK_REASON_BUSY: Set busy: %s\n", __func__, rej->task_uuid); do_remove_uuid = 1; - sais_plat_busy(cb, 1); + sais_plat_busy(sp, 1); break; + case SAI_TASK_REASON_DESTROYED: lwsl_notice("%s: SAI_TASK_REASON_DESTROYED: Clear busy: %s\n", __func__, rej->task_uuid); do_remove_uuid = 1; - sais_plat_busy(cb, 0); + + if (rej->ecode & SAISPRF_EXIT) { + if ((rej->ecode & 0xff) == 0) { + n = SAIES_STEP_SUCCESS; + lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_SUCCESS: %s\n", __func__, rej->task_uuid); + } else { + n = SAIES_FAIL; + lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL: %s\n", __func__, rej->task_uuid); + } + } else + if (rej->ecode & 0x2000) { + n = SAIES_CANCELLED; + lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED: %s\n", __func__, rej->task_uuid); + + } else { + n = SAIES_FAIL; + lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_FAIL: %s\n", __func__, rej->task_uuid); + } + + if (sais_set_task_state(vhd, rej->task_uuid, n, 0, + lws_now_secs() - pss->first_log_timestamp)) + goto bail; + + sais_plat_busy(sp, 0); break; } if (do_remove_uuid && - sais_is_task_inflight(vhd, cb, rej->task_uuid, &ul)) { + sais_is_task_inflight(vhd, sp, rej->task_uuid, &ul)) { lwsl_notice("%s: ### Removing %s from inflight\n", __func__, rej->task_uuid); sais_inflight_entry_destroy(ul); // sais_task_clear_build_and_logs(vhd, rej->task_uuid, 1); @@ -906,11 +878,6 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b /* uuid will not be found listed as inflight for this */ sais_create_and_offer_task_step(vhd, rej->task_uuid, 10); - cb->s_avail_slots = cb->avail_slots; - cb->s_inflight_count = (int)cb->inflight_owner.count; - lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, - sizeof(cb->s_last_rej_task_uuid)); - sais_list_builders(vhd); lwsac_free(&pss->a.ac); @@ -1235,40 +1202,61 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b break; case SAIM_WSSCH_BUILDER_METRIC: - metric = (const sai_build_metric_t *)pss->a.dest; - sais_metrics_db_add(vhd, metric); - - { - uint8_t buf[2048]; - size_t used = 0; - lws_struct_serialize_t *js = lws_struct_json_serialize_create( - lsm_schema_build_metric, - LWS_ARRAY_SIZE(lsm_schema_build_metric), - 0, (void *)metric); - if (js) { - switch (lws_struct_json_serialize(js, buf, sizeof(buf), &used)) { - case LSJS_RESULT_CONTINUE: - assert(0); /* !!! we don't expect to generate anything that won't fit in one fragment */ - break; - case LSJS_RESULT_ERROR: - assert(0); /* we don't expect to not to be able to represent the metrics */ - break; - case LSJS_RESULT_FINISH: - memset(&info, 0, sizeof(info)); + metric = (sai_build_metric_t *)pss->a.dest; + + /* + * We have serialized the incoming JSON representation + * into a sai_build_metric_t *metric. + * + * Let's send it back into JSON so we can broadcast it. + */ - info.private_source_idx = SAI_WEBSRV_PB__PROXIED_FROM_BUILDER; - info.buf = (uint8_t *)buf; - info.len = used; - info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + js = lws_struct_json_serialize_create( + lsm_schema_map_build_metric, + LWS_ARRAY_SIZE(lsm_schema_map_build_metric), + 0, (void *)metric); - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) - lwsl_warn("%s: unable to broadcast to web\n", __func__); + if (!js) + break; - break; - } - lws_struct_json_serialize_destroy(&js); - } + switch (lws_struct_json_serialize(js, xbuf + LWS_PRE, + sizeof(xbuf) - LWS_PRE, &used)) { + case LSJS_RESULT_CONTINUE: + assert(0); /* !!! we don't expect to generate anything that won't fit in one fragment */ + break; + case LSJS_RESULT_ERROR: + assert(0); /* !!! we don't expect to not to be able to represent the metrics */ + break; + case LSJS_RESULT_FINISH: + memset(&info, 0, sizeof(info)); + + info.private_source_idx = SAI_WEBSRV_PB__PROXIED_FROM_BUILDER; + info.buf = xbuf + LWS_PRE; + info.len = used; + info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + + lws_dll2_owner_clear(&o); + lws_dll2_add_head(&metric->list, &o); + + sai_dump_stderr((const char *)xbuf + LWS_PRE, used); + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: unable to broadcast to web\n", __func__); + + /* + * Let's send the struct also into Sqlite3 so we + * can store the metrics + */ + + if (lws_struct_sq3_serialize(pss->vhd->pdb_metrics, + lsm_schema_sq3_map_build_metric, + &o, 0) < 0) + lwsl_err("%s: !!!!!!!!!!!!!!!!!! failed to set metrics in db\n", __func__); + + break; } + lws_struct_json_serialize_destroy(&js); + lwsac_free(&pss->a.ac); break; @@ -1501,10 +1489,7 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsac_free(&task->ac_task_container); free(task); - if ((ssize_t)write(2, start, w) != (ssize_t)w) - lwsl_err("%s: failed to log JSON\n", __func__); - if ((ssize_t)write(2, "\n", 1) != (ssize_t)1) - lwsl_err("%s: failed to log JSON\n", __func__); + sai_dump_stderr((const char *)start, w); lwsl_err("%s: ########## ATTACH TASK --^\n", __func__); diff --git a/src/server/s-ws-web.c b/src/server/s-ws-web.c new file mode 100644 index 0000000..ae4955c --- /dev/null +++ b/src/server/s-ws-web.c @@ -0,0 +1,637 @@ +/* + * Sai server + * + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * b1 --\ sai- sai- /-- browser + * b2 ----- server ---- web ------ browser + * b3 --/ * \-- browser + * + * This is a ws server over a unix domain socket made available by sai-server + * and connected to by sai-web instances running on the same box. + * + * The server notifies the sai-web instances of event and task changes (just + * that a particular event or task changed) and builder list updates (the + * whole current builder list JSON each time). + * + * Sai-web instances can send requests to restart or delete tasks and whole + * events made by authenticated clients. + * + * Since this is on a local UDS protected by user:group, there's no tls or auth + * on this link itself. + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> +#include <assert.h> + +#include "s-private.h" + +typedef struct sai_sul_retry_ctx { + lws_sorted_usec_list_t sul; + struct vhd *vhd; + char uuid[SAI_TASKID_LEN + 1]; + char platform[96]; + int retries; + uint8_t op; /* SAIS_WS_WEBSRV_RX_... */ +} sai_sul_retry_ctx_t; + + +static lws_struct_map_t lsm_browser_taskreset[] = { + LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "uuid"), +}; + +static lws_struct_map_t lsm_browser_platreset[] = { + LSM_CARRAY (sai_browse_rx_platreset_t, event_uuid, "event_uuid"), + LSM_CARRAY (sai_browse_rx_platreset_t, platform, "platform"), +}; + +static const lws_struct_map_t lsm_viewercount_members[] = { + LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"), +}; + +static const lws_struct_map_t lsm_schema_json_map[] = { + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, + /* shares struct */ "com.warmcat.sai.taskreset"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, + /* shares struct */ "com.warmcat.sai.taskrebuildlaststep"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, + /* shares struct */ "com.warmcat.sai.eventreset"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, + /* shares struct */ "com.warmcat.sai.eventdelete"), + LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, + "com.warmcat.sai.taskcan"), + LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewercount_members, + "com.warmcat.sai.viewercount"), + LSM_SCHEMA (sai_rebuild_t, NULL, lsm_rebuild, + "com.warmcat.sai.rebuild"), + LSM_SCHEMA (sai_browse_rx_platreset_t, NULL, lsm_browser_platreset, + "com.warmcat.sai.platreset"), + LSM_SCHEMA (sai_stay_t, NULL, lsm_stay, + "com.warmcat.sai.stay"), +}; + +enum { + SAIS_WS_WEBSRV_RX_TASKRESET, + SAIS_WS_WEBSRV_RX_TASKREBUILDLASTSTEP, + SAIS_WS_WEBSRV_RX_EVENTRESET, + SAIS_WS_WEBSRV_RX_EVENTDELETE, + SAIS_WS_WEBSRV_RX_TASKCANCEL, + SAIS_WS_WEBSRV_RX_VIEWERCOUNT, + SAIS_WS_WEBSRV_RX_REBUILD, + SAIS_WS_WEBSRV_RX_PLATRESET, + SAIS_WS_WEBSRV_RX_STAY, +}; + +static int +sais_validate_id(const char *id, int reqlen) +{ + const char *idin = id; + int n = reqlen; + + while (*id && n--) { + if (!((*id >= '0' && *id <= '9') || + (*id >= 'a' && *id <= 'z') || + (*id >= 'A' && *id <= 'Z'))) + goto reject; + id++; + } + + if (!n && !*id) + return 0; +reject: + + lwsl_notice("%s: Invalid ID (%d) '%s'\n", __func__, reqlen, idin); + + return 1; +} + +static int +sais_validate_builder_name(const char *id) +{ + const char *idin = id; + + while (*id) { + if (!((*id >= '0' && *id <= '9') || + (*id >= 'a' && *id <= 'z') || + (*id >= 'A' && *id <= 'Z') || + *id == '.' || *id == '/' || *id == '-' || *id == '_')) + goto reject; + id++; + } + + return 0; +reject: + + lwsl_notice("%s: Invalid builder name '%s'\n", __func__, idin); + + return 1; +} + + +int +sais_list_builders(struct vhd *vhd) +{ + char json_builders[LWS_PRE + 1024], *start = json_builders + LWS_PRE, + *p = start, *end = p + sizeof(json_builders) - LWS_PRE, + subsequent = 0; + unsigned int ss_flags = LWSSS_FLAG_SOM; + lws_dll2_owner_t db_builders_owner; + sai_plat_t *builder_from_db; + lws_struct_serialize_t *js; + struct lwsac *ac = NULL; + lws_wsmsg_info_t info; + size_t w; + + memset(&db_builders_owner, 0, sizeof(db_builders_owner)); + + if (lws_struct_sq3_deserialize(vhd->server.pdb, NULL, "name ", + lsm_schema_sq3_map_plat, + &db_builders_owner, &ac, 0, 100)) { + lwsl_err("%s: Failed to query builders from DB\n", __func__); + return 1; + } + + // lwsl_warn("%s: count deserialized %d\n", __func__, (int)db_builders_owner.count); + + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), + "{\"schema\":\"com.warmcat.sai.builders\",\"builders\":["); + + lws_start_foreach_dll(struct lws_dll2 *, walk, db_builders_owner.head) { + lws_struct_json_serialize_result_t r; + sai_plat_t *live_builder; + + builder_from_db = lws_container_of(walk, sai_plat_t, sai_plat_list); + live_builder = sais_builder_from_uuid(vhd, builder_from_db->name); + + if (live_builder) { + builder_from_db->online = 1; + lws_strncpy(builder_from_db->peer_ip, live_builder->peer_ip, + sizeof(builder_from_db->peer_ip)); + builder_from_db->stay_on = live_builder->stay_on; + } else + builder_from_db->online = 0; + + builder_from_db->powering_up = 0; + builder_from_db->powering_down = 0; + + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.power_state_owner.head) { + sai_power_state_t *ps = lws_container_of(p, sai_power_state_t, list); + size_t host_len = strlen(ps->host); + + if (!strncmp(builder_from_db->name, ps->host, host_len) && + builder_from_db->name[host_len] == '.') { + builder_from_db->powering_up = ps->powering_up; + builder_from_db->powering_down = ps->powering_down; + break; + } + } lws_end_foreach_dll(p); + + js = lws_struct_json_serialize_create(lsm_schema_map_plat_simple, + LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), + 0, builder_from_db); + if (!js) + goto bail; + + if (subsequent) + *p++ = ','; + subsequent = 1; + + do { + r = lws_struct_json_serialize(js, (uint8_t *)p, + lws_ptr_diff_size_t(end, p) - 2, &w); + p += w; + + switch (r) { + case LSJS_RESULT_FINISH: + /* fallthru */ + case LSJS_RESULT_CONTINUE: + memset(&info, 0, sizeof(info)); + + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)start; + info.len = lws_ptr_diff_size_t(p, start); + info.ss_flags = ss_flags; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: unable to broadcast to web\n", __func__); + + p = start; + ss_flags &= ~((unsigned int)LWSSS_FLAG_SOM); + break; + + case LSJS_RESULT_ERROR: + lws_struct_json_serialize_destroy(&js); + goto bail; + } + + } while (r == LSJS_RESULT_CONTINUE); + + lws_struct_json_serialize_destroy(&js); + + } lws_end_foreach_dll(walk); + + ss_flags |= LWSSS_FLAG_EOM; + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); + memset(&info, 0, sizeof(info)); + + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)start; + info.len = lws_ptr_diff_size_t(p, start); + info.ss_flags = ss_flags; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: unable to broadcast to web\n", __func__); + + // lwsl_notice("%s: Broadcasting builder list: %s\n", __func__, start); + lwsac_free(&ac); + return 0; + +bail: + lwsac_free(&ac); + return 1; +} + + + +static void +sum_viewers_cb(struct lws_ss_handle *h, void *arg) +{ + websrvss_srv_t *m_client = (websrvss_srv_t *)lws_ss_to_user_object(h); + *(unsigned int *)arg += m_client->viewers; +} + + + + +static lws_ss_state_return_t +websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + websrvss_srv_t *m = (websrvss_srv_t *)userobj; + sai_browse_rx_evinfo_t *ei; + lws_struct_args_t a; + sai_db_result_t r; + int n; + + // lwsl_user("%s: len %d, flags: %d\n", __func__, (int)len, flags); + // lwsl_hexdump_info(buf, len); + + memset(&a, 0, sizeof(a)); + a.map_st[0] = lsm_schema_json_map; + a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_json_map); + a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_json_map); + a.ac_block_size = 128; + + lws_struct_json_init_parse(&m->ctx, NULL, &a); + n = lejp_parse(&m->ctx, (uint8_t *)buf, (int)len); + if (n < 0 || !a.dest) { + lwsl_hexdump_notice(buf, len); + lwsl_notice("%s: notification JSON decode failed '%s'\n", + __func__, lejp_error_to_string(n)); + return LWSSSSRET_DISCONNECT_ME; + } + + // lwsl_notice("%s: schema idx %d\n", __func__, a.top_schema_index); + + switch (a.top_schema_index) { + + case SAIS_WS_WEBSRV_RX_TASKRESET: + 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_TASKRESET: %s: received", ei->event_hash); + if (sais_task_clear_build_and_logs(m->vhd, ei->event_hash, 0)) + lwsl_ss_err(m->ss, "taskreset failed"); + 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)) + goto soft_error; + + lwsl_ss_warn(m->ss, "SAIS_WS_WEBSRV_RX_TASKREBUILDLASTSTEP: %s: received", ei->event_hash); + if (sais_task_rebuild_last_step(m->vhd, ei->event_hash)) + lwsl_ss_err(m->ss, "taskrebuildlaststep failed"); + break; + + case SAIS_WS_WEBSRV_RX_EVENTRESET: + ei = (sai_browse_rx_evinfo_t *)a.dest; + + if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) + goto soft_error; + + r = sais_event_reset(m->vhd, ei->event_hash); + if (r) + lwsl_ss_err(m->ss, "eventreset failed"); + + lwsac_free(&a.ac); + break; + + case SAIS_WS_WEBSRV_RX_PLATRESET: { + sai_browse_rx_platreset_t *pr = (sai_browse_rx_platreset_t *)a.dest; + + if (sais_validate_id(pr->event_uuid, SAI_EVENTID_LEN)) + goto soft_error; + + r = sais_plat_reset(m->vhd, pr->event_uuid, pr->platform); + if (r) + lwsl_ss_err(m->ss, "platreset failed"); + lwsac_free(&a.ac); + break; + } + + case SAIS_WS_WEBSRV_RX_EVENTDELETE: + ei = (sai_browse_rx_evinfo_t *)a.dest; + if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) { + lwsl_err("%s: SAIS_WS_WEBSRV_RX_EVENTDELETE: unable to validate id %s\n", __func__, ei->event_hash); + goto soft_error; + } + + lwsl_notice("%s: eventdelete %s\n", __func__, ei->event_hash); + + r = sais_event_delete(m->vhd, ei->event_hash); + if (r) + lwsl_ss_err(m->ss, "event delete failed"); + lwsac_free(&a.ac); + break; + + case SAIS_WS_WEBSRV_RX_TASKCANCEL: + ei = (sai_browse_rx_evinfo_t *)a.dest; + if (sais_validate_id(ei->event_hash, SAI_TASKID_LEN)) + goto soft_error; + + sais_task_cancel(m->vhd, ei->event_hash); + + break; + + case SAIS_WS_WEBSRV_RX_VIEWERCOUNT: + { + sai_viewer_state_t *vs = (sai_viewer_state_t *)a.dest; + unsigned int total_viewers = 0; + char old_viewers_present = !!m->vhd->viewers_are_present; + + /* Store viewer count for this specific sai-web client */ + m->viewers = vs->viewers; + + /* Recalculate total from all connected sai-web clients */ + lws_ss_server_foreach_client(m->vhd->h_ss_websrv, + sum_viewers_cb, &total_viewers); + + m->vhd->browser_viewer_count = total_viewers; + + m->vhd->viewers_are_present = !!total_viewers; + + /* + * Only broadcast to builders if the state has changed + * from 0 viewers to >0, or from >0 viewers to 0. + */ + if (old_viewers_present != m->vhd->viewers_are_present) { + lwsl_notice("%s: Viewer presence changed to %d. Broadcasting to builders.\n", + __func__, m->vhd->viewers_are_present); + lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) { + struct pss *pss_builder = lws_container_of(p, struct pss, same); + sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); + + if (vsend) { + vsend->viewers = m->vhd->viewers_are_present; + lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); + lws_callback_on_writable(pss_builder->wsi); + } + } lws_end_foreach_dll(p); + } + break; + } + case SAIS_WS_WEBSRV_RX_REBUILD: + { + sai_rebuild_t *reb = (sai_rebuild_t *)a.dest; + sai_plat_t *sp; + + if (sais_validate_builder_name(reb->builder_name)) + goto soft_error; + + sp = sais_builder_from_uuid(m->vhd, reb->builder_name); + if (!sp) { + lwsl_info("%s: unknown builder %s for rebuild\n", + __func__, reb->builder_name); + lwsac_free(&a.ac); + break; + } + + /* sp->wsi is the builder connection to server */ + lws_start_foreach_dll(struct lws_dll2 *, p, + m->vhd->builders.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + + if (pss->wsi == sp->wsi) { + sai_rebuild_t *r = malloc(sizeof(*r)); + + if (!r) + break; + *r = *reb; + lws_dll2_add_tail(&r->list, + &pss->rebuild_owner); + lws_callback_on_writable(pss->wsi); + break; + } + } lws_end_foreach_dll(p); + } + break; + + case SAIS_WS_WEBSRV_RX_STAY: + { + sai_stay_t *stay = (sai_stay_t *)a.dest; + + lws_start_foreach_dll(struct lws_dll2 *, p, + m->vhd->sai_powers.head) { + struct pss *pss_power = lws_container_of(p, struct pss, same); + sai_stay_t *s; + + s = malloc(sizeof(*s)); + if (s) { + *s = *stay; + lws_dll2_add_tail(&s->list, &pss_power->stay_owner); + lws_callback_on_writable(pss_power->wsi); + } + } lws_end_foreach_dll(p); + + lwsac_free(&a.ac); + break; + } + } + + return 0; + +soft_error: + lwsl_warn("%s: soft error\n", __func__); + + return 0; +} + +static lws_ss_state_return_t +websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, + size_t *len, int *flags) +{ + websrvss_srv_t *m = (websrvss_srv_t *)userobj; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->bl_srv_to_web), depi, fl; + char som, som1, eom, final = 1; + size_t fsl, used; + + if (!m->bl_srv_to_web) + return LWSSSSRET_TX_DONT_SEND; + + depi = *pi; + + /* + * We can only issue *len at a time. + * + * Notice we are getting the stored flags from the START of the fragment each time. + * that means we can still see the right flags stored with the fragment, even if we + * have partially used the buflist frag and are partway through it. + * + * Ergo, only something to skip if we are at som=1. And also notice that although + * *pi will be right, after the lws_buflist..._use() api, what it points to has been + * destroyed. So we also dereference *pi into depi for use below. + */ + + fsl = lws_buflist_next_segment_len(&m->bl_srv_to_web, NULL); + + lws_buflist_fragment_use(&m->bl_srv_to_web, NULL, 0, &som, &eom); + if (som) { + fsl -= sizeof(int); + lws_buflist_fragment_use(&m->bl_srv_to_web, buf, sizeof(int), &som1, &eom); + } + if (!(depi & LWSSS_FLAG_SOM)) + som = 0; + + used = (size_t)lws_buflist_fragment_use(&m->bl_srv_to_web, (uint8_t *)buf, *len, &som1, &eom); + if (!used) + return LWSSSSRET_TX_DONT_SEND; + + if (used < fsl || !(depi & LWSSS_FLAG_EOM)) /* we saved SS flags at the start of the buf */ + final = 0; + + *len = used; + fl = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); + + // lwsl_ss_notice(m->ss, "Sending %d srv->web: ssflags %d", (int)*len, fl); + + if ((fl & LWSSS_FLAG_SOM) && (((*flags) & 3) == 2)) { + lwsl_ss_err(m->ss, "TX: Illegal LWSSS_FLAG_SOM after previous frame without LWSSS_FLAG_EOM"); + assert(0); + } + if (!(fl & LWSSS_FLAG_SOM) && ((*flags) & 3) == 3) { + lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM after previous frame with LWSSS_FLAG_EOM"); + assert(0); + } + if (!(fl & LWSSS_FLAG_SOM) && !((*flags) & 2)) { + lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM on first frame"); + assert(0); + } + + *flags = fl; + + + // lwsl_hexdump_notice(buf, *len); + + if (m->bl_srv_to_web) + return lws_ss_request_tx(m->ss); + + return 0; +} + + +static lws_ss_state_return_t +websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + websrvss_srv_t *m = (websrvss_srv_t *)userobj; + + // lwsl_user("%s: %p %s, ord 0x%x\n", __func__, m->ss, + // lws_ss_state_name((int)state), (unsigned int)ack); + + switch (state) { + case LWSSSCS_DISCONNECTED: { + unsigned int total_viewers = 0; + + lws_buflist_destroy_all_segments(&m->bl_srv_to_web); + lws_wsmsg_destroy(m->private_heads, LWS_ARRAY_SIZE(m->private_heads)); + + m->viewers = 0; + + /* This sai-web client disconnected, recalculate total viewers */ + lws_ss_server_foreach_client(m->vhd->h_ss_websrv, + sum_viewers_cb, &total_viewers); + + m->vhd->browser_viewer_count = total_viewers; + char new_viewers_present = !!total_viewers; + + if (m->vhd->viewers_are_present != new_viewers_present) { + m->vhd->viewers_are_present = !!new_viewers_present; + lwsl_notice("%s: A sai-web client disconnected, viewer presence changed to %d. Broadcasting.\n", + __func__, new_viewers_present); + + /* Broadcast new presence state to builders */ + lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) { + struct pss *pss_builder = lws_container_of(p, struct pss, same); + sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); + + if (vsend) { + vsend->viewers = (unsigned int)new_viewers_present; + lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); + lws_callback_on_writable(pss_builder->wsi); + } + } lws_end_foreach_dll(p); + } + + break; + } + case LWSSSCS_CREATING: + m->viewers = 0; + return lws_ss_request_tx(m->ss); + + case LWSSSCS_CONNECTED: + sais_list_builders(m->vhd); + break; + case LWSSSCS_ALL_RETRIES_FAILED: + break; + + case LWSSSCS_SERVER_TXN: + break; + + case LWSSSCS_SERVER_UPGRADE: + break; + + default: + break; + } + + return 0; +} + +const lws_ss_info_t ssi_server = { + .handle_offset = offsetof(websrvss_srv_t, ss), + .opaque_user_data_offset = offsetof(websrvss_srv_t, vhd), + .streamtype = "websrv", + .rx = websrvss_ws_rx, + .tx = websrvss_ws_tx, + .state = websrvss_srv_state, + .user_alloc = sizeof(websrvss_srv_t), +}; diff --git a/src/web/CMakeLists.txt b/src/web/CMakeLists.txt index 8f0ea81..481ebf7 100644 --- a/src/web/CMakeLists.txt +++ b/src/web/CMakeLists.txt @@ -6,10 +6,11 @@ set(SRCS w-sai.c w-conf.c w-comms.c - w-central.c w-artifact.c + w-ws-server.c w-ws-browser.c - w-websrv.c + ../common/c-utils.c + ../common/struct-metadata.c ) set(requirements 1) diff --git a/src/web/w-central.c b/src/web/w-central.c deleted file mode 100644 index c94634d..0000000 --- a/src/web/w-central.c +++ /dev/null @@ -1,45 +0,0 @@ -/* - * Sai server - * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - * - * - * Central dispatcher for jobs from events that made it into the database. This - * is done in an event-driven way in m-task.c, but management of it also has to - * be done in the background for when there are no events coming, - */ - -#include <libwebsockets.h> -#include <string.h> -#include <signal.h> -#include <time.h> - -#include "w-private.h" - -extern struct lws_context *context; - -void -saiw_central_cb(lws_sorted_usec_list_t *sul) -{ - struct vhd *vhd = lws_container_of(sul, struct vhd, sul_central); - - /* check again in 1s */ - - lws_sul_schedule(context, 0, &vhd->sul_central, saiw_central_cb, - 1 * LWS_US_PER_SEC); -} diff --git a/src/web/w-comms.c b/src/web/w-comms.c index 2b07fc6..aee5e74 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -1,7 +1,7 @@ /* * Sai web * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -37,29 +37,6 @@ #include "w-private.h" -#include "../common/struct-metadata.c" - -extern const lws_struct_map_t lsm_schema_json_map[]; - -typedef enum { - SJS_CLONING, - SJS_ASSIGNING, - SJS_WAITING, - SJS_DONE -} sai_job_state_t; - -typedef struct sai_job { - struct lws_dll2 jobs_list; - char reponame[64]; - char ref[64]; - char head[64]; - - time_t requested; - - sai_job_state_t state; - -} sai_job_t; - const lws_struct_map_t lsm_schema_map_ta[] = { LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), }; @@ -83,30 +60,6 @@ const lws_struct_map_t lsm_schema_sq3_map_auth[] = { LSM_SCHEMA_DLL2 (sai_auth_t, list, NULL, lsm_auth, "auth"), }; -static lws_struct_map_t lsm_websrv_evinfo[] = { - LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"), -}; - -const lws_struct_map_t lsm_schema_json_map[] = { - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, - /* shares struct */ "sai-taskchange"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, - /* shares struct */ "sai-eventchange"), - LSM_SCHEMA (sai_plat_owner_t, NULL, lsm_plat_list, "com.warmcat.sai.builders"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, - /* shares struct */ "sai-overview"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, - /* shares struct */ "sai-tasklogs"), - LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, - "com.warmcat.sai.loadreport"), - LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, - "com.warmcat.sai.taskactivity"), - LSM_SCHEMA (sai_build_metric_t, NULL, lsm_build_metric, - "com.warmcat.sai.build-metric"), -}; - -size_t lsm_schema_json_map_array_size = LWS_ARRAY_SIZE(lsm_schema_json_map); - extern const lws_struct_map_t lsm_schema_sq3_map_event[]; /* len is typically 16 (event uuid is 32 chars + NUL) @@ -267,7 +220,7 @@ sais_event_db_close(struct vhd *vhd, sqlite3 **ppdb) int sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid) { - char filepath[256], saf[33], r = 0, ra = 0; + char filepath[256], saf[33], ra = 0; lws_strncpy(saf, event_uuid, sizeof(saf)); lws_filename_purify_inplace(saf); @@ -275,27 +228,27 @@ sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid) lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3", vhd->sqlite3_path_lhs, saf); - r = (char)!!unlink(filepath); - if (r) { - lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, filepath, errno); + if (unlink(filepath)) { + lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, + filepath, errno); ra = 1; } lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3-wal", vhd->sqlite3_path_lhs, saf); - r = (char)!!unlink(filepath); - if (r) { - lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, filepath, errno); + if (unlink(filepath)) { + lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, + filepath, errno); ra = 1; } lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3-shm", vhd->sqlite3_path_lhs, saf); - r = (char)!!unlink(filepath); - if (r) { - lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, filepath, errno); + if (unlink(filepath)) { + lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, + filepath, errno); ra = 1; } @@ -306,19 +259,6 @@ sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid) } - -#if 0 -static void -sais_all_browser_on_writable(struct vhd *vhd) -{ - lws_start_foreach_dll(struct lws_dll2 *, mp, vhd->browsers.head) { - struct pss *pss = lws_container_of(mp, struct pss, same); - - lws_callback_on_writable(pss->wsi); - } lws_end_foreach_dll(mp); -} -#endif - typedef enum { SHMUT_NONE = -1, SHMUT_HOOK, @@ -560,9 +500,6 @@ w_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, lwsl_notice("%s: Auth JWK type %d\n", __func__, vhd->jwt_jwk_auth.kty); - lws_sul_schedule(vhd->context, 0, &vhd->sul_central, - saiw_central_cb, 500 * LWS_US_PER_MS); - /* * Reach out to the sai-server part over the SS ws websrv link */ @@ -736,10 +673,6 @@ w_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, return 0; } -// resp = HTTP_STATUS_OK; - -// /* faillthru */ - http_resp: if (lws_add_http_header_status(wsi, (unsigned int)resp, &p, end)) goto bail; @@ -752,8 +685,6 @@ http_resp: case LWS_CALLBACK_HTTP_WRITEABLE: - // lwsl_notice("%s: HTTP_WRITEABLE\n", __func__); - if (!pss || !pss->blob_artifact) break; @@ -875,7 +806,8 @@ http_resp: sr = lws_spa_get_string(pss->spa, EPN_SUCCESS_REDIR); if (!un || !pw || !sr) { - lwsl_notice("%s: missing form args %p %p %p\n",__func__, un, pw, sr); + lwsl_notice("%s: missing form args %p %p %p\n", + __func__, un, pw, sr); pss->spa_failed = 1; goto final; } @@ -1029,39 +961,6 @@ clean_spa: return -1; } -#if 0 - { - const unsigned char *c; - - n = 0; - - do { - int hlen; - - c = lws_token_to_string((enum lws_token_indexes)n); - if (!c) { - n++; - continue; - } - - hlen = lws_hdr_total_length(wsi, (enum lws_token_indexes)n); - if (!hlen || hlen > (int)sizeof(buf) - 1) { - n++; - continue; - } - - if (lws_hdr_copy(wsi, (char *)buf, sizeof buf, (enum lws_token_indexes)n) < 0) - lwsl_wsi_err(wsi, "ESTABLISHED %s (too big)", (char *)c); - else { - buf[sizeof(buf) - 1] = '\0'; - - lwsl_wsi_notice(wsi, "ESTABLISHED %s = %s\n", (char *)c, buf); - } - n++; - } while (c); - } -#endif - /* * What's the situation with a JWT cookie? Normal users won't * have any, but privileged users will have one, and we should diff --git a/src/web/w-private.h b/src/web/w-private.h index 3a8a7a7..64e18fa 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -1,7 +1,7 @@ /* * Sai server definitions src/server/private.h * - * Copyright (C) 2019 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -47,34 +47,6 @@ typedef struct sai_platform { } sai_platform_t; typedef enum { - SAIN_ACTION_INVALID, - SAIN_ACTION_REPO_UPDATED -} sai_notification_action_t; - -typedef struct { - - sai_event_t e; - sai_task_t t; - - char platbuild[4096]; - char platname[96]; - char explicit_platforms[2048]; - - int event_task_index; - - struct lws_b64state b64; - char *saifile; - uint64_t when; - size_t saifile_in_len; - size_t saifile_out_len; - size_t saifile_out_pos; - size_t saifile_in_seen; - sai_notification_action_t action; - - uint8_t nondefault; -} sai_notification_t; - -typedef enum { WSS_IDLE1, WSS_IDLE2, WSS_IDLE3, @@ -136,7 +108,6 @@ struct pss { struct lws_spa *spa; struct lejp_ctx ctx; struct lws_buflist *raw_tx; - sai_notification_t sn; struct lws_dll2 same; /* owner: vhd.browsers */ struct lws_dll2 subs_list; @@ -227,12 +198,11 @@ struct vhd { struct lws_ss_handle *h_ss_websrv; /* client */ - const char *sqlite3_path_lhs; + const char *sqlite3_path_lhs; - lws_dll2_owner_t sqlite3_cache; /* sais_sqlite_cache_t */ - lws_dll2_owner_t tasklog_cache; - lws_sorted_usec_list_t sul_logcache; - lws_sorted_usec_list_t sul_central; /* background task allocation sul */ + lws_dll2_owner_t sqlite3_cache; /* sais_sqlite_cache_t */ + lws_dll2_owner_t tasklog_cache; + lws_sorted_usec_list_t sul_logcache; }; extern struct lws_context * diff --git a/src/web/w-sai.c b/src/web/w-sai.c index cdcfddb..2a7f740 100644 --- a/src/web/w-sai.c +++ b/src/web/w-sai.c @@ -1,7 +1,7 @@ /* * Sai web * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -76,13 +76,13 @@ int main(int argc, const char **argv) logs = atoi(p); lws_set_log_level(logs, NULL); - lwsl_user("Sai Web - Copyright (C) 2019-2020 Andy Green <andy@warmcat.com>\n"); + lwsl_user("Sai Web - Copyright (C) 2019-2025 Andy Green <andy@warmcat.com>\n"); if ((p = lws_cmdline_option(argc, argv, "-c"))) conf = p; context = sai_lws_context_from_json(conf, &info, pprotocols, - default_ss_policy); + default_ss_policy); if (!context) { lwsl_err("lws init failed\n"); return 1; diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c deleted file mode 100644 index 2ec0151..0000000 --- a/src/web/w-websrv.c +++ /dev/null @@ -1,395 +0,0 @@ -/* - * Sai web websrv - saiw SS client private UDS link to sais SS server - * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - * - * We copy JSON to heap and forward it in order to sais side. - */ - -#include <libwebsockets.h> -#include <string.h> -#include <signal.h> -#include <time.h> - -#include "w-private.h" - -typedef struct saiw_websrv { - struct lws_ss_handle *ss; - void *opaque_data; - - lws_struct_args_t a; - struct lejp_ctx ctx; - struct lws_buflist *wbltx; -} saiw_websrv_t; - -extern const lws_struct_map_t lsm_schema_json_map[]; -extern size_t lsm_schema_json_map_array_size; - -enum { - SAIS_WS_WEBSRV_RX_TASKCHANGE, - SAIS_WS_WEBSRV_RX_EVENTCHANGE, - SAIS_WS_WEBSRV_RX_SAI_BUILDERS, - SAIS_WS_WEBSRV_RX_OVERVIEW, /* deleted or added event */ - SAIS_WS_WEBSRV_RX_TASKLOGS, /* new logs for task (ratelimited) */ - SAIS_WS_WEBSRV_RX_LOADREPORT, /* builder's cpu load report */ - SAIS_WS_WEBSRV_RX_TASKACTIVITY, -}; - -/* - * This allows other parts of sai-web to queue a raw buffer to be sent to - * all connected browsers, eg, for load reports. - * - * The flags are lws_write() flags. - */ -void -saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int api_ver_min, enum lws_write_protocol flags) -{ - int eff = 0; - - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { - struct pss *pss = lws_container_of(p, struct pss, same); - int *pi = (int *)((const char *)buf - sizeof(int)); - - eff++; - *pi = (int)flags; - - if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) < 0) - lwsl_wsi_err(pss->wsi, "unable to buflist_append"); /* still ask to drain */ - - lws_callback_on_writable(pss->wsi); - - } lws_end_foreach_dll(p); -} - - -/* - * Queue messages to send from sai-web to sai-server - */ - -int -saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) -{ - saiw_websrv_t *m = (saiw_websrv_t *)lws_ss_to_user_object(h); - unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); - - *pi = ss_flags; - - // lwsl_ss_notice(h, "sai-web: Queuing sai-web -> sai-server"); - // lwsl_hexdump_notice(buf, len); - - if (lws_buflist_append_segment(&m->wbltx, buf - sizeof(int), len + sizeof(int)) < 0) - lwsl_ss_err(h, "failed to append"); /* still ask to drain */ - - if (lws_ss_request_tx(h)) - lwsl_ss_err(h, "failed to request tx"); - - return 0; -} - -/* - * sai-web is receiving from sai-server - * - * This may come in chunks and is statefully parsed - * so it's not directly sensitive to size or fragmentation - */ -static int -saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) -{ - saiw_websrv_t *m = (saiw_websrv_t *)userobj; - struct vhd *vhd = (struct vhd *)m->opaque_data; - sai_browse_rx_evinfo_t *ei; - int n; - - // lwsl_warn("%s: len %d, flags %d\n", __func__, (int)len, flags); - // lwsl_hexdump_notice(buf, len); - - if (flags & LWSSS_FLAG_SOM) { - /* First fragment of a new message. Clear old parse results and init. */ - lwsac_free(&m->a.ac); - memset(&m->a, 0, sizeof(m->a)); - m->a.map_st[0] = lsm_schema_json_map; - m->a.map_entries_st[0] = lsm_schema_json_map_array_size; - m->a.ac_block_size = 4096; - - lws_struct_json_init_parse(&m->ctx, NULL, &m->a); - } - - // fprintf(stderr, "%s: rx: %.*s\n", __func__, (int)len, buf); - - n = lejp_parse(&m->ctx, (uint8_t *)buf, (int)len); - - /* 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; - } - - if (n == LEJP_CONTINUE) { - /* - * Also forward this fragment to browsers if the message is for them. - * We can check the schema index which is available after the - * "schema" member is parsed, even on the first fragment. - */ - switch (m->a.top_schema_index) { - case SAIS_WS_WEBSRV_RX_LOADREPORT: - saiw_ws_broadcast_raw(vhd, buf, len, 0, - lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, 0)); - break; - case SAIS_WS_WEBSRV_RX_TASKACTIVITY: - saiw_ws_broadcast_raw(vhd, buf, len, 0, - lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, 0)); - break; - } - - return 0; - } - - /* - * If we get here, the message is fully parsed (n >= 0). - * Now we can safely process m->a.dest. - */ - if (!m->a.dest) { - lwsl_warn("%s: JSON parsed but produced no object\n", __func__); - goto cleanup_parse_allocs; - } - - switch (m->a.top_schema_index) { - - case SAIS_WS_WEBSRV_RX_TASKCHANGE: - ei = (sai_browse_rx_evinfo_t *)m->a.dest; - 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; - - 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 *cb = lws_container_of(p, sai_plat_t, sai_plat_list); - - lws_dll2_remove(&cb->sai_plat_list); - lws_dll2_add_tail(&cb->sai_plat_list, &vhd->builders_owner); - } lws_end_foreach_dll_safe(p, p1); - - /* 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_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); - } lws_end_foreach_dll(p); - break; - - case SAIS_WS_WEBSRV_RX_OVERVIEW: - lwsl_notice("%s: force overview\n", __func__); - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { - struct pss *pss = lws_container_of(p, struct pss, same); - saiw_alloc_sched(pss, WSS_PREPARE_OVERVIEW); - } lws_end_foreach_dll(p); - break; - - case SAIS_WS_WEBSRV_RX_TASKLOGS: - ei = (sai_browse_rx_evinfo_t *)m->a.dest; - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->subs_owner.head) { - struct pss *pss = lws_container_of(p, struct pss, subs_list); - if (!strcmp(pss->sub_task_uuid, ei->event_hash)) - lws_callback_on_writable(pss->wsi); - } lws_end_foreach_dll(p); - break; - - case SAIS_WS_WEBSRV_RX_LOADREPORT: - // lwsl_notice("%s: ^^^^^^^^^^^^^^ SAIS_WS_WEBSRV_RX_LOADREPORT forwarding to browser\n", __func__); - // lwsl_hexdump_notice(buf, len); - saiw_ws_broadcast_raw(vhd, buf, len - (unsigned int)n, 0, - 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_raw(vhd, buf, len - (unsigned int)n, 0, - lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); - break; - } - -cleanup_parse_allocs: - /* - * Free the memory used for THIS parse. - * In the BUILDERS case, m->a.ac was transferred to vhd->builders, - * so it will be NULL here and lwsac_free is a no-op. - */ - lwsac_free(&m->a.ac); - return 0; - -cleanup_and_disconnect: - lwsac_free(&m->a.ac); - return LWSSSSRET_DISCONNECT_ME; -} - - -static int -saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, - int *flags) -{ - saiw_websrv_t *m = (saiw_websrv_t *)userobj; - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->wbltx), depi; - char som, som1, eom, final = 1; - size_t fsl, used; - - if (!m->wbltx) { - lwsl_notice("%s: nothing to send from web -> srv\n", __func__); - return LWSSSSRET_TX_DONT_SEND; - } - - depi = *pi; - - /* - * We can only issue *len at a time. - * - * Notice we are getting the stored flags from the START of the fragment each time. - * that means we can still see the right flags stored with the fragment, even if we - * have partially used the buflist frag and are partway through it. - * - * Ergo, only something to skip if we are at som=1. And also notice that although - * *pi will be right, after the lws_buflist..._use() api, what it points to has been - * destroyed. So we also dereference *pi into depi for use below. - */ - - fsl = lws_buflist_next_segment_len(&m->wbltx, NULL); - - lws_buflist_fragment_use(&m->wbltx, NULL, 0, &som, &eom); - if (som) { - fsl -= sizeof(int); - lws_buflist_fragment_use(&m->wbltx, buf, sizeof(int), &som1, &eom); - } - - /* this is the only buflist user on pss->raw_tx */ - used = (size_t)lws_buflist_fragment_use(&m->wbltx, (uint8_t *)buf, *len, &som1, &eom); - if (!used) - return LWSSSSRET_TX_DONT_SEND; - - if (used < fsl || (depi & LWS_WRITE_NO_FIN)) - final = 0; - - *len = used; - *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - - // lwsl_ss_notice(m->ss, "Sending %d web->srv: ssflags %d", (int)*len, (int)*flags); - // lwsl_hexdump_notice(buf, *len); - - if (m->wbltx) - return lws_ss_request_tx(m->ss); - - return 0; -} - -static int -saiw_lp_state(void *userobj, void *sh, lws_ss_constate_t state, - lws_ss_tx_ordinal_t ack) -{ - saiw_websrv_t *m = (saiw_websrv_t *)userobj; - struct vhd *vhd = (struct vhd *)m->opaque_data; - - lwsl_info("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name((int)state), - (unsigned int)ack); - - switch (state) { - case LWSSSCS_DESTROYING: - break; - - case LWSSSCS_CONNECTED: - lwsl_info("%s: connected to websrv uds\n", __func__); - return lws_ss_request_tx(m->ss); - - case LWSSSCS_DISCONNECTED: - lws_buflist_destroy_all_segments(&m->wbltx); - lwsac_detach(&vhd->builders); - break; - - case LWSSSCS_ALL_RETRIES_FAILED: - return lws_ss_client_connect(m->ss); - - case LWSSSCS_QOS_ACK_REMOTE: - break; - - default: - break; - } - - return 0; -} - -const lws_ss_info_t ssi_saiw_websrv = { - .handle_offset = offsetof(saiw_websrv_t, ss), - .opaque_user_data_offset = offsetof(saiw_websrv_t, opaque_data), - .rx = saiw_lp_rx, - .tx = saiw_lp_tx, - .state = saiw_lp_state, - .user_alloc = sizeof(saiw_websrv_t), - .streamtype = "websrv" -}; - -/* - * This function calculates the current number of connected browsers and - * sends an update to the sai-server. - */ -void -saiw_update_viewer_count(struct vhd *vhd) -{ - sai_viewer_state_t vs; - char buf[LWS_PRE + 256]; - size_t len; - - if (!vhd || !vhd->h_ss_websrv) - return; - - /* The count is simply the number of items in the browsers list */ - vs.viewers = (unsigned int)vhd->browsers.count; - - const lws_struct_map_t lsm_viewercount_members[] = { - LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"), - }; - - const lws_struct_map_t lsm_schema_json_map[] = { - LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewercount_members, - "com.warmcat.sai.viewercount"), - }; - - lws_struct_serialize_t *js = lws_struct_json_serialize_create( - lsm_schema_json_map, LWS_ARRAY_SIZE(lsm_schema_json_map), - 0, &vs); - if (!js) - return; - - len = 0; - lws_struct_json_serialize(js, (unsigned char *)buf + LWS_PRE, - sizeof(buf) - LWS_PRE, &len); - lws_struct_json_serialize_destroy(&js); - - if (len > 0) - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf + LWS_PRE, len, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); -} diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index d4cd44b..63bc32d 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -18,6 +18,10 @@ * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, * MA 02110-1301 USA * + * b1 --\ sai- sai- /-- browser + * b2 ----- server ---- web ------ browser + * b3 --/ * \-- browser + * * These are ws rx and tx handlers related to browser ws connections, on * /broswe URLs. */ @@ -302,16 +306,12 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) n = lws_struct_sq3_deserialize(pdb, qu, NULL, lsm_schema_sq3_map_task, &o, &sch->query_ac, 0, 1); sais_event_db_close(pss->vhd, &pdb); - lwsl_notice("%s: WWWWWWWWWWW -- actual task n %d, o.head %p\n", __func__, n, o.head); if (n < 0 || !o.head) goto bail; pt = lws_container_of(o.head, sai_task_t, list); sch->one_task = pt; - lwsl_notice("%s: WWWWWWWWWWW -- browser ws asked for task hash: %s, plat %s\n", - __func__, task_uuid, sch->one_task->platform); - /* let the pss take over the task info ac and schedule sending */ lws_dll2_remove((struct lws_dll2 *)&sch->one_task->list); @@ -348,8 +348,6 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) n = lws_struct_sq3_deserialize(pss->vhd->pdb, qu, NULL, lsm_schema_sq3_map_event, &o, &sch->query_ac, 0, 1); - lwsl_notice("%s: WWWWWWWWWWW -- actual event n %d, o.head %p\n", __func__, n, o.head); - if (n < 0 || !o.head) /* * It's OK if the parent event is not visible in the current @@ -1192,14 +1190,8 @@ so_finish: pss->send_state = WSS_IDLE1; saiw_dealloc_sched(sch); return 1; - case LSJS_RESULT_FINISH: - p += w; - lws_struct_json_serialize_destroy(&js); - sch->walk = sch->walk->next; - if (!sch->walk) - goto b_finish; - break; + case LSJS_RESULT_FINISH: case LSJS_RESULT_CONTINUE: p += w; lws_struct_json_serialize_destroy(&js); @@ -1212,6 +1204,7 @@ so_finish: break; b_finish: p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); + // lwsac_unreference(&vhd->builders); endo = 1; break; @@ -1314,10 +1307,8 @@ b_finish: } lwsl_notice("%s: wwwwwwwwwww TASKINFO\n", __func__); - if ((size_t)write(2, start, lws_ptr_diff_size_t(p, start)) != lws_ptr_diff_size_t(p, start)) - lwsl_notice("%s: dump JSON failed\n", __func__); - lwsl_notice("\n"); + sai_dump_stderr((const char *)start, lws_ptr_diff_size_t(p, start)); break; case WSS_SEND_ARTIFACT_INFO: diff --git a/src/web/w-ws-server.c b/src/web/w-ws-server.c new file mode 100644 index 0000000..20527de --- /dev/null +++ b/src/web/w-ws-server.c @@ -0,0 +1,445 @@ +/* + * Sai web websrv - saiw SS client private link to sais SS server + * + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * b1 --\ sai- sai- /-- browser + * b2 ----- server ---- web ------ browser + * b3 --/ * \-- browser + * + * We copy JSON to heap and forward it in order to sais side. + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> + +#include "w-private.h" + +typedef struct saiw_websrv { + struct lws_ss_handle *ss; + void *opaque_data; + + lws_struct_args_t a; + struct lejp_ctx ctx; + struct lws_buflist *wbltx; +} saiw_websrv_t; + +static lws_struct_map_t lsm_websrv_evinfo[] = { + LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"), +}; + +const lws_struct_map_t lsm_schema_json_map[] = { + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, + /* shares struct */ "sai-taskchange"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, + /* shares struct */ "sai-eventchange"), + LSM_SCHEMA (sai_plat_owner_t, NULL, lsm_plat_list, "com.warmcat.sai.builders"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, + /* shares struct */ "sai-overview"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, + /* shares struct */ "sai-tasklogs"), + LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, + "com.warmcat.sai.loadreport"), + LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, + "com.warmcat.sai.taskactivity"), + LSM_SCHEMA (sai_build_metric_t, NULL, lsm_build_metric, + "com.warmcat.sai.build-metric"), +}; + +enum { + SAIS_WS_WEBSRV_RX_TASKCHANGE, + SAIS_WS_WEBSRV_RX_EVENTCHANGE, + SAIS_WS_WEBSRV_RX_SAI_BUILDERS, + SAIS_WS_WEBSRV_RX_OVERVIEW, /* deleted or added event */ + SAIS_WS_WEBSRV_RX_TASKLOGS, /* new logs for task (ratelimited) */ + SAIS_WS_WEBSRV_RX_LOADREPORT, /* builder's cpu load report */ + SAIS_WS_WEBSRV_RX_TASKACTIVITY, + SAIS_WS_WEBSRV_RX_BUILD_METRIC, +}; + +/* + * This allows other parts of sai-web to queue a raw buffer to be sent to + * all connected browsers, eg, for load reports. + * + * The flags are lws_write() flags. + */ +void +saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int api_ver_min, enum lws_write_protocol flags) +{ + int eff = 0; + + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + int *pi = (int *)((const char *)buf - sizeof(int)); + + eff++; + *pi = (int)flags; + + if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) < 0) + lwsl_wsi_err(pss->wsi, "unable to buflist_append"); /* still ask to drain */ + + lws_callback_on_writable(pss->wsi); + + } lws_end_foreach_dll(p); +} + + +/* + * Queue messages to send from sai-web to sai-server + */ + +int +saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) +{ + saiw_websrv_t *m = (saiw_websrv_t *)lws_ss_to_user_object(h); + unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); + + *pi = ss_flags; + + // lwsl_ss_notice(h, "sai-web: Queuing sai-web -> sai-server"); + // lwsl_hexdump_notice(buf, len); + + if (lws_buflist_append_segment(&m->wbltx, buf - sizeof(int), len + sizeof(int)) < 0) + lwsl_ss_err(h, "failed to append"); /* still ask to drain */ + + if (lws_ss_request_tx(h)) + lwsl_ss_err(h, "failed to request tx"); + + return 0; +} + +/* + * sai-web is receiving from sai-server + * + * This may come in chunks and is statefully parsed + * so it's not directly sensitive to size or fragmentation + */ +static int +saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + saiw_websrv_t *m = (saiw_websrv_t *)userobj; + struct vhd *vhd = (struct vhd *)m->opaque_data; + sai_browse_rx_evinfo_t *ei; + int n; + + // 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) { + /* First frag of a new message. Clear old parse results and init */ + lwsac_free(&m->a.ac); + memset(&m->a, 0, sizeof(m->a)); + m->a.map_st[0] = lsm_schema_json_map; + m->a.map_entries_st[0] = 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); + } + + // fprintf(stderr, "%s: rx: %.*s\n", __func__, (int)len, buf); + + 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; + } + + if (n == LEJP_CONTINUE) { + /* + * Also forward this fragment to browsers if the message is for them. + * We can check the schema index which is available after the + * "schema" member is parsed, even on the first fragment. + */ + switch (m->a.top_schema_index) { + case SAIS_WS_WEBSRV_RX_LOADREPORT: + saiw_ws_broadcast_raw(vhd, buf, len, 0, + lws_write_ws_flags(LWS_WRITE_TEXT, + flags & LWSSS_FLAG_SOM, 0)); + break; + case SAIS_WS_WEBSRV_RX_TASKACTIVITY: + saiw_ws_broadcast_raw(vhd, buf, len, 0, + lws_write_ws_flags(LWS_WRITE_TEXT, + flags & LWSSS_FLAG_SOM, 0)); + break; + case SAIS_WS_WEBSRV_RX_SAI_BUILDERS: + saiw_ws_broadcast_raw(vhd, buf, len, 0, + 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; + } + + return 0; + } else { + 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: + saiw_ws_broadcast_raw(vhd, buf, len, 0, + lws_write_ws_flags(LWS_WRITE_TEXT, + flags & LWSSS_FLAG_SOM, + flags & LWSSS_FLAG_EOM)); + break; + } + // lwsl_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; + + 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_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); + } lws_end_foreach_dll(p); + break; + + case SAIS_WS_WEBSRV_RX_OVERVIEW: + lwsl_notice("%s: force overview\n", __func__); + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + + saiw_alloc_sched(pss, WSS_PREPARE_OVERVIEW); + } lws_end_foreach_dll(p); + break; + + case SAIS_WS_WEBSRV_RX_TASKLOGS: + ei = (sai_browse_rx_evinfo_t *)m->a.dest; + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->subs_owner.head) { + struct pss *pss = lws_container_of(p, struct pss, subs_list); + if (!strcmp(pss->sub_task_uuid, ei->event_hash)) + lws_callback_on_writable(pss->wsi); + } lws_end_foreach_dll(p); + break; + + case SAIS_WS_WEBSRV_RX_LOADREPORT: + // lwsl_notice("%s: ^^^^^^^^^^^^^^ SAIS_WS_WEBSRV_RX_LOADREPORT forwarding to browser\n", __func__); + // lwsl_hexdump_notice(buf, len); + saiw_ws_broadcast_raw(vhd, buf, len - (unsigned int)n, 0, + 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_raw(vhd, buf, len - (unsigned int)n, 0, + lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); + break; + } + +cleanup_parse_allocs: + /* + * Free the memory used for THIS parse. + * In the BUILDERS case, m->a.ac was transferred to vhd->builders, + * so it will be NULL here and lwsac_free is a no-op. + */ + lwsac_free(&m->a.ac); + return 0; + +cleanup_and_disconnect: + lwsac_free(&m->a.ac); + return LWSSSSRET_DISCONNECT_ME; +} + + +static int +saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, + int *flags) +{ + saiw_websrv_t *m = (saiw_websrv_t *)userobj; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->wbltx), depi; + char som, som1, eom, final = 1; + size_t fsl, used; + + if (!m->wbltx) { + // lwsl_notice("%s: nothing to send from web -> srv\n", __func__); + return LWSSSSRET_TX_DONT_SEND; + } + + depi = *pi; + + /* + * We can only issue *len at a time. + * + * Notice we are getting the stored flags from the START of the fragment each time. + * that means we can still see the right flags stored with the fragment, even if we + * have partially used the buflist frag and are partway through it. + * + * Ergo, only something to skip if we are at som=1. And also notice that although + * *pi will be right, after the lws_buflist..._use() api, what it points to has been + * destroyed. So we also dereference *pi into depi for use below. + */ + + fsl = lws_buflist_next_segment_len(&m->wbltx, NULL); + + lws_buflist_fragment_use(&m->wbltx, NULL, 0, &som, &eom); + if (som) { + fsl -= sizeof(int); + lws_buflist_fragment_use(&m->wbltx, buf, sizeof(int), &som1, &eom); + } + + /* this is the only buflist user on pss->raw_tx */ + used = (size_t)lws_buflist_fragment_use(&m->wbltx, (uint8_t *)buf, *len, &som1, &eom); + if (!used) + return LWSSSSRET_TX_DONT_SEND; + + if (used < fsl || (depi & LWS_WRITE_NO_FIN)) + final = 0; + + *len = used; + *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); + + // lwsl_ss_notice(m->ss, "Sending %d web->srv: ssflags %d", (int)*len, (int)*flags); + // lwsl_hexdump_notice(buf, *len); + + if (m->wbltx) + return lws_ss_request_tx(m->ss); + + return 0; +} + +static int +saiw_lp_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + saiw_websrv_t *m = (saiw_websrv_t *)userobj; + struct vhd *vhd = (struct vhd *)m->opaque_data; + + lwsl_info("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name((int)state), + (unsigned int)ack); + + switch (state) { + case LWSSSCS_DESTROYING: + break; + + case LWSSSCS_CONNECTED: + lwsl_info("%s: connected to websrv uds\n", __func__); + return lws_ss_request_tx(m->ss); + + case LWSSSCS_DISCONNECTED: + lws_buflist_destroy_all_segments(&m->wbltx); + lwsac_detach(&vhd->builders); + break; + + case LWSSSCS_ALL_RETRIES_FAILED: + return lws_ss_client_connect(m->ss); + + case LWSSSCS_QOS_ACK_REMOTE: + break; + + default: + break; + } + + return 0; +} + +const lws_ss_info_t ssi_saiw_websrv = { + .handle_offset = offsetof(saiw_websrv_t, ss), + .opaque_user_data_offset = offsetof(saiw_websrv_t, opaque_data), + .rx = saiw_lp_rx, + .tx = saiw_lp_tx, + .state = saiw_lp_state, + .user_alloc = sizeof(saiw_websrv_t), + .streamtype = "websrv" +}; + +/* + * This function calculates the current number of connected browsers and + * sends an update to the sai-server. + */ +void +saiw_update_viewer_count(struct vhd *vhd) +{ + sai_viewer_state_t vs; + char buf[LWS_PRE + 256]; + size_t len; + + if (!vhd || !vhd->h_ss_websrv) + return; + + /* The count is simply the number of items in the browsers list */ + vs.viewers = (unsigned int)vhd->browsers.count; + + const lws_struct_map_t lsm_viewercount_members[] = { + LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"), + }; + + const lws_struct_map_t lsm_schema_json_map[] = { + LSM_SCHEMA (sai_viewer_state_t, NULL, lsm_viewercount_members, + "com.warmcat.sai.viewercount"), + }; + + lws_struct_serialize_t *js = lws_struct_json_serialize_create( + lsm_schema_json_map, LWS_ARRAY_SIZE(lsm_schema_json_map), + 0, &vs); + if (!js) + return; + + len = 0; + lws_struct_json_serialize(js, (unsigned char *)buf + LWS_PRE, + sizeof(buf) - LWS_PRE, &len); + lws_struct_json_serialize_destroy(&js); + + if (len > 0) + saiw_websrv_queue_tx(vhd->h_ss_websrv, buf + LWS_PRE, len, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); +}
Page fetched 0s ago, creation time: 29ms (vhost etag hits: 0%, cache hits: 0%)