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 / arch-aarch64-a72-bcm2711-rpi4.svg
Author[]Andy Green <andy@warmcat.com> 2025-10-19 12:29 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-10-22 12:10 UTC
Tree4868e559bc2b5e8e2b214bf9cb9fe4cfe592fe83   Raw Patch
 
clean: improve frags from srv to web
clean: improve frags from srv to web
diff --git a/assets/sai.css b/assets/sai.css index a186e51..d1d2470 100644 --- a/assets/sai.css +++ b/assets/sai.css @@ -582,6 +582,10 @@ td.summary { padding: 4px; } +td.builder-info { + vertical-align: top; +} + td.bn { font-weight: normal; font-size: 7pt; @@ -628,7 +632,7 @@ div.ib { div.ibuil { display: inline-block; - vertical-align:middle; + vertical-align:top; text-align:left; line-height: 50%; float:left; diff --git a/assets/sai.js b/assets/sai.js index 82d241e..4575bae 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -800,10 +800,10 @@ function summarize_build_situation(event_uuid) text = total + " pending"; else { var parts = []; - if (good) parts.push(good + " OK"); - if (bad) parts.push(bad + " bad"); - if (ongoing) parts.push(ongoing + " building"); - if (pending) parts.push(pending + " wait"); + if (good) parts.push("OK: " + good); + if (bad) parts.push("Bad: " + bad); + if (ongoing) parts.push("Building: " + ongoing); + if (pending) parts.push("Wait: " + pending); text = parts.join(", "); } diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index 9df46ba..3707bae 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -74,7 +74,7 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, { lws_struct_serialize_t *js; unsigned int ssf = LWSSS_FLAG_SOM; - uint8_t buf[4096]; + uint8_t buf[1024]; size_t w = 0; js = lws_struct_json_serialize_create(map, map_entries, 0, object); @@ -86,10 +86,8 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, do { switch (lws_struct_json_serialize(js, buf, sizeof(buf), &w)) { case LSJS_RESULT_CONTINUE: - lwsl_notice("%s: LSJS_RESULT_CONTINUE\n", __func__); break; case LSJS_RESULT_FINISH: - lwsl_notice("%s: LSJS_RESULT_FINISH\n", __func__); ssf |= LWSSS_FLAG_EOM; break; case LSJS_RESULT_ERROR: @@ -97,9 +95,6 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, return -1; } - lwsl_notice("%s: queueing %d bytes, ss_flags %d\n", __func__, (int)w, ssf); - lwsl_hexdump_notice(buf, w); - if (saib_srv_queue_tx(h, buf, w, ssf)) return -1; @@ -117,8 +112,8 @@ saib_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) struct sai_plat_server *spm = (struct sai_plat_server *)userobj; //struct sai_plat *sp = (struct sai_plat *)spm->sai_plat; - lwsl_info("%s: len %d, flags: %d\n", __func__, (int)len, flags); - lwsl_hexdump_info(buf, len); +// lwsl_info("%s: len %d, flags: %d\n", __func__, (int)len, flags); +// lwsl_hexdump_info(buf, len); if (saib_ws_json_rx_builder(spm, buf, len)) return 1; @@ -147,6 +142,7 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, } depi = *pi; + *pi = (*pi) & (~(LWSSS_FLAG_SOM)); /* no SOM twice even on partial */ /* * We can only issue *len at a time. @@ -167,19 +163,21 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, 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 & LWS_WRITE_NO_FIN)) + 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 web->srv: ssflags %d", (int)*len, (int)*flags); - lwsl_hexdump_notice(buf, *len); +// 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); diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index 0121f3c..781a682 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -372,10 +372,15 @@ skip: free(op); } + if (ns->task) + saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid, + SAI_TASK_REASON_DESTROYED); + return; fail: - n = lws_snprintf(s, sizeof(s), "Build step %d FAILED, exit code: %d\n", ns->current_step + 1, exit_code); + n = lws_snprintf(s, sizeof(s), "Build step %d FAILED, exit code: %d\n", + ns->current_step + 1, exit_code); saib_log_chunk_create(ns, s, (size_t)n, 3); saib_task_grace(ns); @@ -383,6 +388,10 @@ fail: saib_log_chunk_create(ns, NULL, 0, 2); + if (ns->task) + saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid, + SAI_TASK_REASON_DESTROYED); + if (op->spawn) free(op->spawn); diff --git a/src/builder/b-private.h b/src/builder/b-private.h index 445e318..b02aa51 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -267,3 +267,6 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, const lws_struct_map_t *map, size_t map_entries, void *object); +int +saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, + const char *rej_task_uuid, unsigned int reason); diff --git a/src/builder/b-refproxy.c b/src/builder/b-refproxy.c index 30f2dae..cd9f37e 100644 --- a/src/builder/b-refproxy.c +++ b/src/builder/b-refproxy.c @@ -103,8 +103,8 @@ saib_handle_resource_result(struct sai_plat_server *spm, const char *in, size_t pss = resproxy_find_by_cookie(spm, p, al); if (!pss) { - lwsl_warn("%s: the requestor left before the response\n", - __func__); + // lwsl_warn("%s: the requestor left before the response\n", + // __func__); /* * Explicit yield, in case the acceptance raced the client * closing... if it was telling us we can't have it, the server diff --git a/src/builder/b-task.c b/src/builder/b-task.c index 81a2843..e72bab5 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -281,7 +281,7 @@ saib_set_ns_state(struct sai_nspawn *ns, int state) * update all servers we're connected to about builder status / optional reject */ -static int +int saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, const char *rej_task_uuid, unsigned int reason) { @@ -368,18 +368,14 @@ saib_task_destroy(struct sai_nspawn *ns) m++; } lws_end_foreach_dll_safe(d, d1); - if (!m) - lws_sul_schedule(builder.context, 0, - &builder.sul_idle, sul_idle_cb, - SAI_IDLE_GRACE_US); - /* * Schedule informing all the servers we're connected to */ - if (ns->task) - saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid, - SAI_TASK_REASON_DESTROYED); + if (!m) + lws_sul_schedule(builder.context, 0, + &builder.sul_idle, sul_idle_cb, + SAI_IDLE_GRACE_US); } if (ns->task && ns->task->ac_task_container) { diff --git a/src/common/include/private.h b/src/common/include/private.h index 79cce7d..8da9761 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -437,6 +437,7 @@ typedef struct sai_plat_server_ref { /* common struct for lists of task uuids on a builder */ typedef struct sai_uuid_list { lws_dll2_t list; + lws_usec_t us_time_listed; char uuid[65]; char started; } sai_uuid_list_t; @@ -633,8 +634,10 @@ extern const lws_struct_map_t lsm_schema_build_metric[1], lsm_schema_sq3_map_build_metric[1], lsm_load_report_members[9], - lsm_schema_json_task_rej[5] - ; + 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]; @@ -658,3 +661,5 @@ sul_idle_cb(lws_sorted_usec_list_t *sul); int sai_uuid16_create(struct lws_context *context, char *dest33); + + diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index 774e0f7..49e50a0 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -7,10 +7,14 @@ set(SRCS s-conf.c s-notification.c s-comms.c + s-helpers.c + s-power.c s-ws-builder.c s-task.c + s-task-helpers.c s-central.c s-websrv.c + s-webops.c s-resource.c s-metrics-db.c ../common/c-utils.c diff --git a/src/server/s-comms.c b/src/server/s-comms.c index fc4fdb4..13c4f5a 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -1,7 +1,7 @@ /* * Sai server * - * 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 @@ -61,286 +61,6 @@ typedef struct sai_job { extern const lws_struct_map_t lsm_schema_sq3_map_event[]; extern const lws_ss_info_t ssi_server; -/* len is typically 16 (event uuid is 32 chars + NUL) - * But eg, task uuid is concatenated 32-char eventid and 32-char taskid - */ - -int -sai_sqlite3_statement(sqlite3 *pdb, const char *cmd, const char *desc) -{ - sqlite3_stmt *sm; - int n; - - if (sqlite3_prepare_v2(pdb, cmd, -1, &sm, NULL) != SQLITE_OK) { - lwsl_err("%s: Unable to %s: %s\n", - __func__, desc, sqlite3_errmsg(pdb)); - - return 1; - } - - n = sqlite3_step(sm); - sqlite3_reset(sm); - sqlite3_finalize(sm); - if (n != SQLITE_DONE) { - n = sqlite3_extended_errcode(pdb); - if (!n) { - lwsl_info("%s: failed '%s'\n", __func__, cmd); - return 0; - } - - lwsl_err("%s: %d: Unable to perform \"%s\": %s\n", __func__, - n, desc, sqlite3_errmsg(pdb)); - puts(cmd); - - return 1; - } - - return 0; -} - -int -sais_event_db_ensure_open(struct vhd *vhd, const char *event_uuid, - char create_if_needed, sqlite3 **ppdb) -{ - char filepath[256], saf[33]; - sais_sqlite_cache_t *sc; - - // lwsl_notice("%s: (sai-server) entry\n", __func__); - - if (*ppdb) - return 0; - - /* do we have this guy cached? */ - - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->sqlite3_cache.head) { - sc = lws_container_of(p, sais_sqlite_cache_t, list); - - if (!strcmp(event_uuid, sc->uuid)) { - sc->refcount++; - *ppdb = sc->pdb; - return 0; - } - - } lws_end_foreach_dll(p); - - /* ... nope, well, let's open and cache him then... */ - - lws_strncpy(saf, event_uuid, sizeof(saf)); - lws_filename_purify_inplace(saf); - - lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3", - vhd->sqlite3_path_lhs, saf); - - if (lws_struct_sq3_open(vhd->context, filepath, create_if_needed, ppdb)) { - lwsl_err("%s: Unable to open db %s: %s\n", __func__, - filepath, sqlite3_errmsg(*ppdb)); - - return 2; - } - - /* create / add to the schema for the tables we will have in here */ - - if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_task)) { - lwsl_err("%s: unable to create task table in %s\n", __func__, filepath); - return 3; - } - - sai_sqlite3_statement(*ppdb, "PRAGMA journal_mode=WAL;", "set WAL"); - - if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_log)) { - lwsl_err("%s: unable to create log table in %s\n", __func__, filepath); - - return 4; - } - - if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_artifact)) { - lwsl_err("%s: unable to create artifact table in %s\n", __func__, filepath); - - return 5; - } - - sc = malloc(sizeof(*sc)); - memset(sc, 0, sizeof(*sc)); - if (!sc) { - lwsl_err("%s: unable to alloc sc for %s\n", __func__, filepath); - - lws_struct_sq3_close(ppdb); - *ppdb = NULL; - return 6; - } - - lws_strncpy(sc->uuid, event_uuid, sizeof(sc->uuid)); - sc->refcount = 1; - sc->pdb = *ppdb; - lws_dll2_add_tail(&sc->list, &vhd->sqlite3_cache); - - return 0; -} - -void -sais_event_db_close(struct vhd *vhd, sqlite3 **ppdb) -{ - sais_sqlite_cache_t *sc; - - if (!*ppdb) - return; - - /* look for him in the cache */ - - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->sqlite3_cache.head) { - sc = lws_container_of(p, sais_sqlite_cache_t, list); - - if (sc->pdb == *ppdb) { - *ppdb = NULL; - if (--sc->refcount) { - lwsl_notice("%s: zero refcount to idle\n", - __func__); - /* - * He's not currently in use then... don't - * close him immediately, s-central.c has a - * timer that closes and removes sqlite3 - * cache entries idle for longer than 60s - */ - sc->idle_since = lws_now_usecs(); - } - - return; - } - - } lws_end_foreach_dll(p); - - lws_struct_sq3_close(ppdb); - *ppdb = NULL; -} - -int -sais_event_db_close_all_now(struct vhd *vhd) -{ - sais_sqlite_cache_t *sc; - - lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, - vhd->sqlite3_cache.head) { - sc = lws_container_of(p, sais_sqlite_cache_t, list); - - lws_struct_sq3_close(&sc->pdb); - lws_dll2_remove(&sc->list); - free(sc); - - } lws_end_foreach_dll_safe(p, p1); - - return 0; -} - -int -sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid) -{ - char filepath[256], saf[33], r = 0, ra = 0; - - lws_strncpy(saf, event_uuid, sizeof(saf)); - lws_filename_purify_inplace(saf); - - lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3", - vhd->sqlite3_path_lhs, saf); - - r = (char)!!unlink(filepath); - if (r) { - lwsl_err("%s: 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: 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: unable to delete %s (%d)\n", __func__, filepath, errno); - ra = 1; - } - - if (!ra) - lwsl_notice("%s: deleted %s OK\n", __func__, filepath); - - return ra; -} - - -#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 - -static int -sai_detach_builder(struct lws_dll2 *d, void *user) -{ -// saib_t *b = lws_container_of(d, saib_t, c.builder_list); - - lws_dll2_remove(d); - - return 0; -} - -static int -sai_detach_resource(struct lws_dll2 *d, void *user) -{ - lws_dll2_remove(d); - - return 0; -} - -static int -sai_destroy_resource_wellknown(struct lws_dll2 *d, void *user) -{ - sai_resource_wellknown_t *rwk = - lws_container_of(d, sai_resource_wellknown_t, list); - - /* - * Just detach everything listed on this well-known resource... - * everything listed here is ultimately owned by a pss and will be - * destroyed when that goes down - */ - - lws_dll2_foreach_safe(&rwk->owner_queued, NULL, sai_detach_resource); - lws_dll2_foreach_safe(&rwk->owner_leased, NULL, sai_detach_resource); - - lws_dll2_remove(d); - - free(rwk); - - return 0; -} - -static void -sais_server_destroy(struct vhd *vhd, sais_t *server) -{ - lwsl_notice("%s: server %p\n", __func__, server); - if (server) - lws_dll2_foreach_safe(&server->builder_owner, NULL, - sai_detach_builder); - - sais_event_db_close_all_now(vhd); - - lws_struct_sq3_close(&server->pdb); - - lws_dll2_foreach_safe(&server->resource_wellknown_owner, NULL, - sai_destroy_resource_wellknown); -} - typedef enum { SHMUT_NONE = -1, SHMUT_HOOK, @@ -392,6 +112,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, struct pss *pss = (struct pss *)user; sai_http_murl_t mu = SHMUT_NONE; const char *pvo_resources, *num; + unsigned int ssf; int n; (void)end; @@ -672,9 +393,12 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, * they can update connected browsers to show the new * event */ - - sais_websrv_broadcast(vhd->h_ss_websrv, - "{\"schema\":\"sai-overview\"}", 25); + n = lws_snprintf((char *)start, sizeof(buf) - LWS_PRE, + "{\"schema\":\"sai-overview\"}"); + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, + (const char *)start, (size_t)n, + SAI_WEBSRV_PB__GENERATED, + LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); } if (lws_return_http_status(wsi, @@ -773,21 +497,6 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, pss->pdb_artifact = NULL; } - /* drop any inflight task information for this builder */ - - lws_start_foreach_dll(struct lws_dll2 *, pb, - vhd->server.builder_owner.head) { - sai_plat_t *build = lws_container_of(pb, sai_plat_t, sai_plat_list); - - lws_start_foreach_dll_safe(struct lws_dll2 *, pif, pif1, - build->inflight_owner.head) { - sai_uuid_list_t *ul = lws_container_of(pif, sai_uuid_list_t, list); - - sais_inflight_entry_destroy(ul); - - } lws_end_foreach_dll_safe(pif, pif1); - } lws_end_foreach_dll(pb); - /* * Update the sai-webs about the builder removal, so they * can update their connected browsers @@ -798,6 +507,10 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, case LWS_CALLBACK_RECEIVE: + pss->wsi = wsi; + ssf = (lws_is_first_fragment(wsi) ? LWSSS_FLAG_SOM : 0) | + (lws_is_final_fragment(wsi) ? LWSSS_FLAG_EOM : 0); + /* * 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` @@ -805,128 +518,17 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, */ if (pss->is_power) { - struct lejp_ctx ctx; - lws_struct_args_t a; - sai_power_state_t *ps; - const lws_struct_map_t lsm_schema_map_power[] = { - LSM_SCHEMA(sai_power_state_t, NULL, lsm_power_state, - "com.warmcat.sai.powerstate"), - LSM_SCHEMA(sai_power_managed_builders_t, NULL, - lsm_power_managed_builders_list, - "com.warmcat.sai.power_managed_builders"), - LSM_SCHEMA(sai_stay_state_update_t, NULL, - lsm_stay_state_update, - "com.warmcat.sai.stay_state_update"), - }; - - /* This is a message from sai-power */ - lwsl_notice("RX from sai-power: %.*s\n", (int)len, (const char *)in); - - memset(&a, 0, sizeof(a)); - a.map_st[0] = lsm_schema_map_power; - a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_power); - a.ac_block_size = 512; - - lws_struct_json_init_parse(&ctx, NULL, &a); - if (lejp_parse(&ctx, (uint8_t *)in, (int)len) < 0 || !a.dest) { - lwsl_warn("Failed to parse msg from sai-power\n"); - lwsac_free(&a.ac); - break; // Exit case - } - - switch (a.top_schema_index) { - case 0: /* powerstate */ - ps = (sai_power_state_t *)a.dest; - if (ps->powering_up) { - lwsl_notice("sai-power is powering up: %s\n", ps->host); - sais_set_builder_power_state(vhd, ps->host, 1, 0); - } else if (ps->powering_down) { - lwsl_notice("sai-power is powering down: %s\n", ps->host); - sais_set_builder_power_state(vhd, ps->host, 0, 1); - } - break; - - case 1: { - sai_power_managed_builders_t *pmb = (sai_power_managed_builders_t *)a.dest; - uint64_t bf_set = 0; - - lws_start_foreach_dll(struct lws_dll2 *, p, pmb->builders.head) { - sai_power_managed_builder_t *b = lws_container_of(p, - sai_power_managed_builder_t, list); - char q[256]; - int shi = 0; - - lwsl_notice("%s: Marking builder %s as power-managed\n", - __func__, b->name); - lws_snprintf(q, sizeof(q), - "UPDATE builders SET power_managed=1 WHERE name = '%s' OR name LIKE '%s.%%'", - b->name, b->name); - if (sai_sqlite3_statement(vhd->server.pdb, q, "set power_managed")) - lwsl_err("%s: Failed to mark builder %s as power-managed\n", - __func__, b->name); - - 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, '.'); - - // 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; - bf_set |= (1 << shi); - } // else - // lwsl_notice("%s: ------------ Unmatched '%s' '%s'\n", __func__, cb->name + (dot - cb->name) - strlen(b->name), b->name); - - shi++; - } lws_end_foreach_dll(p2); - - } lws_end_foreach_dll(p); - - sais_list_builders(vhd); - - break; - } - case 2: { - sai_stay_state_update_t *ssu = (sai_stay_state_update_t *)a.dest; - sai_plat_t *cb; - - 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, - sai_plat_list); - - const char *dot = strchr(cb->name, '.'); - - if (dot && !strncmp(cb->name, ssu->builder_name, (size_t)(dot - cb->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; - sais_list_builders(vhd); - break; - } - } lws_end_foreach_dll(p); - - break; - } - } - - lwsac_free(&a.ac); + sais_power_rx(vhd, pss, in, len, ssf); break; } /* * This is a message from a builder */ - // lwsl_notice("%s: rx from builder, len %d, final: %d\n", __func__, (int)len, lws_is_final_fragment(wsi)); - pss->wsi = wsi; - if (sais_ws_json_rx_builder(vhd, pss, in, len)) + + // lwsl_wsi_notice(wsi, "rx from builder, len %d, : ss_flags: %d\n", (int)len, ssf); + + if (sais_ws_json_rx_builder(vhd, pss, in, len, ssf)) return -1; if (!pss->announced) { diff --git a/src/server/s-helpers.c b/src/server/s-helpers.c new file mode 100644 index 0000000..a605bab --- /dev/null +++ b/src/server/s-helpers.c @@ -0,0 +1,339 @@ +/* + * 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 + * + * The same ws interface is connected-to by builders (on path /builder), and + * provides the query transport for browsers (on path /browse). + * + * There's a single server slite3 database containing events, and a separate + * sqlite3 database file for each event, it only contains tasks and logs for + * the event and can be deleted when the event record associated with it is + * deleted. This is to keep is scalable when there may be thousands of events + * and related tasks and logs stored. + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> +#include <stdio.h> +#include <fcntl.h> + +#include "s-private.h" + +int +sql3_get_integer_cb(void *user, int cols, char **values, char **name) +{ + unsigned int *pui = (unsigned int *)user; + + if (cols < 1 || !values[0]) + *pui = 0; + else + *pui = (unsigned int)atoi(values[0]); + + return 0; +} + +int +sql3_get_string_cb(void *user, int cols, char **values, char **name) +{ + char *p = (char *)user; + + p[0] = '\0'; + if (cols < 1 || !values[0]) + return 0; + + lws_strncpy(p, values[0], 33); + + return 0; +} + +void +sai_task_uuid_to_event_uuid(char *event_uuid33, const char *task_uuid65) +{ + memcpy(event_uuid33, task_uuid65, 32); + event_uuid33[32] = '\0'; +} + +/* len is typically 16 (event uuid is 32 chars + NUL) + * But eg, task uuid is concatenated 32-char eventid and 32-char taskid + */ + +int +sai_sqlite3_statement(sqlite3 *pdb, const char *cmd, const char *desc) +{ + sqlite3_stmt *sm; + int n; + + if (sqlite3_prepare_v2(pdb, cmd, -1, &sm, NULL) != SQLITE_OK) { + lwsl_err("%s: Unable to %s: %s\n", + __func__, desc, sqlite3_errmsg(pdb)); + + return 1; + } + + n = sqlite3_step(sm); + sqlite3_reset(sm); + sqlite3_finalize(sm); + if (n != SQLITE_DONE) { + n = sqlite3_extended_errcode(pdb); + if (!n) { + lwsl_info("%s: failed '%s'\n", __func__, cmd); + return 0; + } + + lwsl_err("%s: %d: Unable to perform \"%s\": %s\n", __func__, + n, desc, sqlite3_errmsg(pdb)); + puts(cmd); + + return 1; + } + + return 0; +} + +int +sais_event_db_ensure_open(struct vhd *vhd, const char *event_uuid, + char create_if_needed, sqlite3 **ppdb) +{ + char filepath[256], saf[33]; + sais_sqlite_cache_t *sc; + + // lwsl_notice("%s: (sai-server) entry\n", __func__); + + if (*ppdb) + return 0; + + /* do we have this guy cached? */ + + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->sqlite3_cache.head) { + sc = lws_container_of(p, sais_sqlite_cache_t, list); + + if (!strcmp(event_uuid, sc->uuid)) { + sc->refcount++; + *ppdb = sc->pdb; + return 0; + } + + } lws_end_foreach_dll(p); + + /* ... nope, well, let's open and cache him then... */ + + lws_strncpy(saf, event_uuid, sizeof(saf)); + lws_filename_purify_inplace(saf); + + lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3", + vhd->sqlite3_path_lhs, saf); + + if (lws_struct_sq3_open(vhd->context, filepath, create_if_needed, ppdb)) { + lwsl_err("%s: Unable to open db %s: %s\n", __func__, + filepath, sqlite3_errmsg(*ppdb)); + + return 2; + } + + /* create / add to the schema for the tables we will have in here */ + + if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_task)) { + lwsl_err("%s: unable to create task table in %s\n", __func__, filepath); + return 3; + } + + sai_sqlite3_statement(*ppdb, "PRAGMA journal_mode=WAL;", "set WAL"); + + if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_log)) { + lwsl_err("%s: unable to create log table in %s\n", __func__, filepath); + + return 4; + } + + if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_artifact)) { + lwsl_err("%s: unable to create artifact table in %s\n", __func__, filepath); + + return 5; + } + + sc = malloc(sizeof(*sc)); + memset(sc, 0, sizeof(*sc)); + if (!sc) { + lwsl_err("%s: unable to alloc sc for %s\n", __func__, filepath); + + lws_struct_sq3_close(ppdb); + *ppdb = NULL; + return 6; + } + + lws_strncpy(sc->uuid, event_uuid, sizeof(sc->uuid)); + sc->refcount = 1; + sc->pdb = *ppdb; + lws_dll2_add_tail(&sc->list, &vhd->sqlite3_cache); + + return 0; +} + +void +sais_event_db_close(struct vhd *vhd, sqlite3 **ppdb) +{ + sais_sqlite_cache_t *sc; + + if (!*ppdb) + return; + + /* look for him in the cache */ + + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->sqlite3_cache.head) { + sc = lws_container_of(p, sais_sqlite_cache_t, list); + + if (sc->pdb == *ppdb) { + *ppdb = NULL; + if (--sc->refcount) { + lwsl_notice("%s: zero refcount to idle\n", + __func__); + /* + * He's not currently in use then... don't + * close him immediately, s-central.c has a + * timer that closes and removes sqlite3 + * cache entries idle for longer than 60s + */ + sc->idle_since = lws_now_usecs(); + } + + return; + } + + } lws_end_foreach_dll(p); + + lws_struct_sq3_close(ppdb); + *ppdb = NULL; +} + +int +sais_event_db_close_all_now(struct vhd *vhd) +{ + sais_sqlite_cache_t *sc; + + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + vhd->sqlite3_cache.head) { + sc = lws_container_of(p, sais_sqlite_cache_t, list); + + lws_struct_sq3_close(&sc->pdb); + lws_dll2_remove(&sc->list); + free(sc); + + } lws_end_foreach_dll_safe(p, p1); + + return 0; +} + +int +sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid) +{ + char filepath[256], saf[33], r = 0, ra = 0; + + lws_strncpy(saf, event_uuid, sizeof(saf)); + lws_filename_purify_inplace(saf); + + lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3", + vhd->sqlite3_path_lhs, saf); + + r = (char)!!unlink(filepath); + if (r) { + lwsl_err("%s: 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: 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: unable to delete %s (%d)\n", __func__, filepath, errno); + ra = 1; + } + + if (!ra) + lwsl_notice("%s: deleted %s OK\n", __func__, filepath); + + return ra; +} + + +int +sai_detach_builder(struct lws_dll2 *d, void *user) +{ + lws_dll2_remove(d); + + return 0; +} + +int +sai_detach_resource(struct lws_dll2 *d, void *user) +{ + lws_dll2_remove(d); + + return 0; +} + +int +sai_destroy_resource_wellknown(struct lws_dll2 *d, void *user) +{ + sai_resource_wellknown_t *rwk = + lws_container_of(d, sai_resource_wellknown_t, list); + + /* + * Just detach everything listed on this well-known resource... + * everything listed here is ultimately owned by a pss and will be + * destroyed when that goes down + */ + + lws_dll2_foreach_safe(&rwk->owner_queued, NULL, sai_detach_resource); + lws_dll2_foreach_safe(&rwk->owner_leased, NULL, sai_detach_resource); + + lws_dll2_remove(d); + + free(rwk); + + return 0; +} + +void +sais_server_destroy(struct vhd *vhd, sais_t *server) +{ + lwsl_notice("%s: server %p\n", __func__, server); + if (server) + lws_dll2_foreach_safe(&server->builder_owner, NULL, + sai_detach_builder); + + sais_event_db_close_all_now(vhd); + + lws_struct_sq3_close(&server->pdb); + + lws_dll2_foreach_safe(&server->resource_wellknown_owner, NULL, + sai_destroy_resource_wellknown); +} + diff --git a/src/server/s-power.c b/src/server/s-power.c new file mode 100644 index 0000000..30153e0 --- /dev/null +++ b/src/server/s-power.c @@ -0,0 +1,158 @@ +/* + * 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 + * + * The same ws interface is connected-to by builders (on path /builder), and + * provides the query transport for browsers (on path /browse). + * + * There's a single server slite3 database containing events, and a separate + * sqlite3 database file for each event, it only contains tasks and logs for + * the event and can be deleted when the event record associated with it is + * deleted. This is to keep is scalable when there may be thousands of events + * and related tasks and logs stored. + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> +#include <stdio.h> +#include <fcntl.h> + +#include "s-private.h" + +int +sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, + size_t bl, unsigned int ss_flags) +{ + struct lejp_ctx ctx; + lws_struct_args_t a; + sai_power_state_t *ps; + const lws_struct_map_t lsm_schema_map_power[] = { + LSM_SCHEMA(sai_power_state_t, NULL, lsm_power_state, + "com.warmcat.sai.powerstate"), + LSM_SCHEMA(sai_power_managed_builders_t, NULL, + lsm_power_managed_builders_list, + "com.warmcat.sai.power_managed_builders"), + LSM_SCHEMA(sai_stay_state_update_t, NULL, + lsm_stay_state_update, + "com.warmcat.sai.stay_state_update"), + }; + + /* This is a message from sai-power */ + lwsl_notice("RX from sai-power: %.*s\n", (int)bl, (const char *)buf); + + memset(&a, 0, sizeof(a)); + a.map_st[0] = lsm_schema_map_power; + a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_power); + a.ac_block_size = 512; + + lws_struct_json_init_parse(&ctx, NULL, &a); + if (lejp_parse(&ctx, buf, (int)bl) < 0 || !a.dest) { + lwsl_warn("Failed to parse msg from sai-power\n"); + lwsac_free(&a.ac); + return 1; + } + + switch (a.top_schema_index) { + case 0: /* powerstate */ + ps = (sai_power_state_t *)a.dest; + if (ps->powering_up) { + lwsl_notice("sai-power is powering up: %s\n", ps->host); + sais_set_builder_power_state(vhd, ps->host, 1, 0); + } else if (ps->powering_down) { + lwsl_notice("sai-power is powering down: %s\n", ps->host); + sais_set_builder_power_state(vhd, ps->host, 0, 1); + } + break; + + case 1: { + sai_power_managed_builders_t *pmb = (sai_power_managed_builders_t *)a.dest; + uint64_t bf_set = 0; + + lws_start_foreach_dll(struct lws_dll2 *, p, pmb->builders.head) { + sai_power_managed_builder_t *b = lws_container_of(p, + sai_power_managed_builder_t, list); + char q[256]; + int shi = 0; + + lwsl_notice("%s: Marking builder %s as power-managed\n", + __func__, b->name); + lws_snprintf(q, sizeof(q), + "UPDATE builders SET power_managed=1 WHERE name = '%s' OR name LIKE '%s.%%'", + b->name, b->name); + + if (sai_sqlite3_statement(vhd->server.pdb, q, "set power_managed")) + lwsl_err("%s: Failed to mark builder %s as power-managed\n", + __func__, b->name); + + 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, '.'); + + // 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; + bf_set |= (1 << shi); + } + shi++; + } lws_end_foreach_dll(p2); + + } lws_end_foreach_dll(p); + + sais_list_builders(vhd); + + break; + } + case 2: { + sai_stay_state_update_t *ssu = (sai_stay_state_update_t *)a.dest; + sai_plat_t *cb; + + 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, + sai_plat_list); + + const char *dot = strchr(cb->name, '.'); + + if (dot && !strncmp(cb->name, ssu->builder_name, (size_t)(dot - cb->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; + sais_list_builders(vhd); + break; + } + } lws_end_foreach_dll(p); + + break; + } + } + + 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 7acb40a..f2c187e 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -28,6 +28,18 @@ struct sai_plat; +/* lws_wsmsg_ array for different sources */ +enum { + SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, + SAI_WEBSRV_PB__LOGS, + SAI_WEBSRV_PB__GENERATED, + SAI_WEBSRV_PB__ACTIVITY, + + SAI_WEBSRV_PB__COUNT +}; + + + typedef enum { SAI_DB_RESULT_OK, SAI_DB_RESULT_BUSY, @@ -57,6 +69,17 @@ typedef struct sai_platform { /* build and name over-allocated here */ } sai_platform_t; +typedef struct websrvss_srv { + struct lws_ss_handle *ss; + struct vhd *vhd; + + struct lejp_ctx ctx; + struct lws_buflist *bl_srv_to_web; + unsigned int viewers; + + struct lws_buflist *private_heads[SAI_WEBSRV_PB__COUNT]; +} websrvss_srv_t; + typedef struct sai_powering_up_plat { lws_dll2_t list; char name[256]; @@ -189,8 +212,6 @@ struct vhd { struct lws_ss_handle *h_ss_websrv; /* server */ - char json_builders[8192]; - /* pss lists */ struct lws_dll2_owner builders; struct lws_dll2_owner sai_powers; @@ -256,7 +277,7 @@ int saiw_ws_json_tx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl); int -sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl); +sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl, unsigned int ss_flags); int sais_list_builders(struct vhd *vhd); @@ -305,8 +326,8 @@ 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); -void -sais_websrv_broadcast(struct lws_ss_handle *hsrv, const char *str, size_t len); +int +sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv, const char *str, size_t len, int reassembly_idx, unsigned int ss_flags); int sql3_get_integer_cb(void *user, int cols, char **values, char **name); @@ -358,4 +379,45 @@ sais_add_to_inflight_list_if_absent(struct vhd *vhd, sai_plat_t *sp, const char void sais_inflight_entry_destroy(sai_uuid_list_t *ul); +void +sais_prune_inflight_list(struct vhd *vhd); + +sai_db_result_t +sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform); + +sai_db_result_t +sais_event_delete(struct vhd *vhd, const char *event_uuid); +sai_db_result_t +sais_event_reset(struct vhd *vhd, const char *event_uuid); + +int +sai_detach_builder(struct lws_dll2 *d, void *user); + +int +sai_detach_resource(struct lws_dll2 *d, void *user); + +int +sai_destroy_resource_wellknown(struct lws_dll2 *d, void *user); + +void +sais_server_destroy(struct vhd *vhd, sais_t *server); + +void +sais_get_task_metrics_estimates(struct vhd *vhd, sai_task_t *task); + +int +sais_task_cancel(struct vhd *vhd, const char *task_uuid); + +int +sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid); + +sai_db_result_t +sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_rejection); + +sai_db_result_t +sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid); + +int +sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, + size_t bl, unsigned int ss_flags); diff --git a/src/server/s-task-helpers.c b/src/server/s-task-helpers.c new file mode 100644 index 0000000..553d08a --- /dev/null +++ b/src/server/s-task-helpers.c @@ -0,0 +1,326 @@ +/* + * 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 + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> +#include <assert.h> + +#include "s-private.h" + +void +sais_get_task_metrics_estimates(struct vhd *vhd, sai_task_t *task) +{ + char query[256]; + sqlite3_stmt *stmt; + + task->est_peak_mem_kib = 256 * 1024; /* 256MiB default */ + task->est_cpu_load_pct = 10; + task->est_disk_kib = 1024 * 1024; /* 1GiB default */ + + if (!vhd->pdb_metrics) + return; + + 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'", + task->taskname); + + 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); + + 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_disk_kib = (unsigned int)(sqlite3_column_int(stmt, 2) / 1024); + } + + sqlite3_finalize(stmt); +} + +int +sais_task_cancel(struct vhd *vhd, const char *task_uuid) +{ + sai_cancel_t *can; + + /* + * For every pss that we have from builders... + */ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->builders.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + + + /* + * ... queue the task cancel message + */ + can = malloc(sizeof *can); + if (!can) + return -1; + memset(can, 0, sizeof(*can)); + + lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); + + lws_dll2_add_tail(&can->list, &pss->task_cancel_owner); + + lws_callback_on_writable(pss->wsi); + + } lws_end_foreach_dll(p); + + sais_taskchange(vhd->h_ss_websrv, task_uuid, SAIES_CANCELLED); + + /* + * Recompute startable task platforms and broadcast to all sai-power, + * after there has been a change in tasks + */ + sais_platforms_with_tasks_pending(vhd); + + return 0; +} + +int +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; + + lwsl_notice("%s: builders count %d\n", __func__, vhd->builders.count); + + /* + * We will send the task cancel message only to the builder that was + * assigned the task, if any. + */ + + sai_task_uuid_to_event_uuid(event_uuid, task_uuid); + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + lwsl_err("%s: unable to open event-specific database\n", __func__); + return -1; + } + + builder_name[0] = '\0'; + lws_sql_purify(esc_uuid, 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]) { + sais_event_db_close(vhd, &pdb); + /* + * This is not an error... the task may not have had a builder + * assigned yet. There's nothing to do. + */ + return 0; + } + sais_event_db_close(vhd, &pdb); + + cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); + if (!cb) + /* 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) { + pss_match = pss; + break; + } + } lws_end_foreach_dll(p); + + if (!pss_match) + /* Builder is live but has no pss? */ + return 0; + + can = malloc(sizeof *can); + if (!can) + return -1; + + memset(can, 0, sizeof(*can)); + + lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); + + lws_dll2_add_tail(&can->list, &pss_match->task_cancel_owner); + lws_callback_on_writable(pss_match->wsi); + + return 0; +} + +/* + * Keep the task record itself, but remove all logs and artifacts related to + * it and reset the task state back to WAITING. + */ + +sai_db_result_t +sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_rejection) +{ + char esc[96], cmd[256], event_uuid[33]; + sqlite3 *pdb = NULL; + int ret; + + 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)) { + lwsl_err("%s: unable to open event-specific database\n", + __func__); + + return SAI_DB_RESULT_ERROR; + } + + lws_sql_purify(esc, task_uuid, sizeof(esc)); + lws_snprintf(cmd, sizeof(cmd), "delete from logs where task_uuid='%s'", + esc); + + ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + lwsl_err("%s: %s: %s: fail\n", __func__, cmd, + sqlite3_errmsg(pdb)); + return SAI_DB_RESULT_ERROR; + } + lws_snprintf(cmd, sizeof(cmd), "delete from artifacts where task_uuid='%s'", + esc); + + ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + lwsl_err("%s: %s: %s: fail\n", __func__, cmd, + sqlite3_errmsg(pdb)); + return SAI_DB_RESULT_ERROR; + } + + sais_event_db_close(vhd, &pdb); + + sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 1, 1); + + sais_task_stop_on_builders(vhd, task_uuid); + + /* + * Reassess now if there's a builder we can match to a pending task, + * but not if we are being reset due to a rejection... that would + * just cause us to spam the builder with the same task again + */ + + if (!from_rejection) { + lwsl_err("%s: scheduling sul_central to find a new task\n", __func__); + lws_sul_schedule(vhd->context, 0, &vhd->sul_central, sais_central_cb, 1); + } + + /* + * Recompute startable task platforms and broadcast to all sai-power, + * after there has been a change in tasks + */ + sais_platforms_with_tasks_pending(vhd); + + lwsl_notice("%s: exiting OK\n", __func__); + + return SAI_DB_RESULT_OK; +} + +sai_db_result_t +sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) +{ + char esc[96], cmd[256], event_uuid[33]; + sqlite3 *pdb = NULL; + lws_dll2_owner_t o; + struct lwsac *ac = NULL; + sai_task_t *task; + int ret; + + if (!task_uuid[0]) + return SAI_DB_RESULT_OK; + + lwsl_notice("%s: received request to rebuild last step of 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)) { + lwsl_err("%s: unable to open event-specific database\n", + __func__); + + return SAI_DB_RESULT_ERROR; + } + + lws_sql_purify(esc, task_uuid, sizeof(esc)); + lws_snprintf(cmd, sizeof(cmd), " and uuid='%s'", esc); + ret = lws_struct_sq3_deserialize(pdb, cmd, NULL, + lsm_schema_sq3_map_task, &o, &ac, 0, 1); + if (ret < 0 || !o.head) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + return SAI_DB_RESULT_ERROR; + } + + task = lws_container_of(o.head, sai_task_t, list); + + if (task->build_step > 0) { + lws_snprintf(cmd, sizeof(cmd), + "update tasks set build_step=%d where uuid='%s'", + task->build_step - 1, esc); + + ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + + lwsl_err("%s: %s: %s: fail\n", __func__, cmd, + sqlite3_errmsg(pdb)); + return SAI_DB_RESULT_ERROR; + } + } + + lwsac_free(&ac); + sais_event_db_close(vhd, &pdb); + + sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 0, 0); + + sais_task_stop_on_builders(vhd, task_uuid); + + lwsl_err("%s: scheduling sul_central to find a new task\n", __func__); + lws_sul_schedule(vhd->context, 0, &vhd->sul_central, sais_central_cb, 1); + + sais_platforms_with_tasks_pending(vhd); + + lwsl_notice("%s: exiting OK\n", __func__); + + return SAI_DB_RESULT_OK; +} \ No newline at end of file diff --git a/src/server/s-task.c b/src/server/s-task.c index bef4371..133882d 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -1,7 +1,7 @@ /* * Sai server * - * 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 @@ -27,33 +27,6 @@ #include "s-private.h" -int -sql3_get_integer_cb(void *user, int cols, char **values, char **name) -{ - unsigned int *pui = (unsigned int *)user; - - if (cols < 1 || !values[0]) - *pui = 0; - else - *pui = (unsigned int)atoi(values[0]); - - return 0; -} - -int -sql3_get_string_cb(void *user, int cols, char **values, char **name) -{ - char *p = (char *)user; - - p[0] = '\0'; - if (cols < 1 || !values[0]) - return 0; - - lws_strncpy(p, values[0], 33); - - return 0; -} - /* temporary info about a task that failed in a previous run */ typedef struct sai_failed_task_info { lws_dll2_t list; @@ -62,13 +35,6 @@ typedef struct sai_failed_task_info { const char *taskname; } sai_failed_task_info_t; -void -sai_task_uuid_to_event_uuid(char *event_uuid33, const char *task_uuid65) -{ - memcpy(event_uuid33, task_uuid65, 32); - event_uuid33[32] = '\0'; -} - 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, @@ -395,6 +361,7 @@ sais_add_to_inflight_list_if_absent(struct vhd *vhd, sai_plat_t *sp, const char memset(uuid_list, 0, sizeof(*uuid_list)); lws_strncpy(uuid_list->uuid, uuid, sizeof(uuid_list->uuid)); + uuid_list->us_time_listed = lws_now_usecs(); lws_dll2_add_tail(&uuid_list->list, &sp->inflight_owner); @@ -412,6 +379,26 @@ sais_inflight_entry_destroy(sai_uuid_list_t *ul) free(ul); } +void +sais_prune_inflight_list(struct vhd *vhd) +{ + lws_usec_t t = lws_now_usecs(); + + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { + sai_plat_t *sp = lws_container_of(p, sai_plat_t, sai_plat_list); + + lws_start_foreach_dll_safe(struct lws_dll2 *, p1, p2, sp->inflight_owner.head) { + sai_uuid_list_t *u = lws_container_of(p1, sai_uuid_list_t, list); + + if (!u->started && (t - u->us_time_listed) > 5 * 1000 * 1000) + sais_inflight_entry_destroy(u); + + } lws_end_foreach_dll_safe(p1, p2); + + } lws_end_foreach_dll(p); +} + + /* * Find the most recent task that still needs doing for platform, on any event */ @@ -489,6 +476,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, do { sqlite3_stmt *sm; + prev_event_uuid[0] = '\0'; lws_snprintf(query, sizeof(query), "select uuid, created from events where repo_name='%s' and " @@ -810,304 +798,6 @@ bail: return 1; } -static void -sais_get_task_metrics_estimates(struct vhd *vhd, sai_task_t *task) -{ - char query[256]; - sqlite3_stmt *stmt; - - task->est_peak_mem_kib = 256 * 1024; /* 256MiB default */ - task->est_cpu_load_pct = 10; - task->est_disk_kib = 1024 * 1024; /* 1GiB default */ - - if (!vhd->pdb_metrics) - return; - - 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'", - task->taskname); - - 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); - - 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_disk_kib = (unsigned int)(sqlite3_column_int(stmt, 2) / 1024); - } - - sqlite3_finalize(stmt); -} - -int -sais_task_cancel(struct vhd *vhd, const char *task_uuid) -{ - sai_cancel_t *can; - - /* - * For every pss that we have from builders... - */ - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->builders.head) { - struct pss *pss = lws_container_of(p, struct pss, same); - - - /* - * ... queue the task cancel message - */ - can = malloc(sizeof *can); - if (!can) - return -1; - memset(can, 0, sizeof(*can)); - - lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); - - lws_dll2_add_tail(&can->list, &pss->task_cancel_owner); - - lws_callback_on_writable(pss->wsi); - - } lws_end_foreach_dll(p); - - sais_taskchange(vhd->h_ss_websrv, task_uuid, SAIES_CANCELLED); - - /* - * Recompute startable task platforms and broadcast to all sai-power, - * after there has been a change in tasks - */ - sais_platforms_with_tasks_pending(vhd); - - return 0; -} - -static int -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; - - lwsl_notice("%s: builders count %d\n", __func__, vhd->builders.count); - - /* - * We will send the task cancel message only to the builder that was - * assigned the task, if any. - */ - - sai_task_uuid_to_event_uuid(event_uuid, task_uuid); - - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { - lwsl_err("%s: unable to open event-specific database\n", __func__); - return -1; - } - - builder_name[0] = '\0'; - lws_sql_purify(esc_uuid, 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]) { - sais_event_db_close(vhd, &pdb); - /* - * This is not an error... the task may not have had a builder - * assigned yet. There's nothing to do. - */ - return 0; - } - sais_event_db_close(vhd, &pdb); - - cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); - if (!cb) - /* 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) { - pss_match = pss; - break; - } - } lws_end_foreach_dll(p); - - if (!pss_match) - /* Builder is live but has no pss? */ - return 0; - - can = malloc(sizeof *can); - if (!can) - return -1; - - memset(can, 0, sizeof(*can)); - - lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); - - lws_dll2_add_tail(&can->list, &pss_match->task_cancel_owner); - lws_callback_on_writable(pss_match->wsi); - - return 0; -} - -/* - * Keep the task record itself, but remove all logs and artifacts related to - * it and reset the task state back to WAITING. - */ - -sai_db_result_t -sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_rejection) -{ - char esc[96], cmd[256], event_uuid[33]; - sqlite3 *pdb = NULL; - int ret; - - 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)) { - lwsl_err("%s: unable to open event-specific database\n", - __func__); - - return SAI_DB_RESULT_ERROR; - } - - lws_sql_purify(esc, task_uuid, sizeof(esc)); - lws_snprintf(cmd, sizeof(cmd), "delete from logs where task_uuid='%s'", - esc); - - ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - lwsl_err("%s: %s: %s: fail\n", __func__, cmd, - sqlite3_errmsg(pdb)); - return SAI_DB_RESULT_ERROR; - } - lws_snprintf(cmd, sizeof(cmd), "delete from artifacts where task_uuid='%s'", - esc); - - ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - lwsl_err("%s: %s: %s: fail\n", __func__, cmd, - sqlite3_errmsg(pdb)); - return SAI_DB_RESULT_ERROR; - } - - sais_event_db_close(vhd, &pdb); - - sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 1, 1); - - sais_task_stop_on_builders(vhd, task_uuid); - - /* - * Reassess now if there's a builder we can match to a pending task, - * but not if we are being reset due to a rejection... that would - * just cause us to spam the builder with the same task again - */ - - if (!from_rejection) { - lwsl_err("%s: scheduling sul_central to find a new task\n", __func__); - lws_sul_schedule(vhd->context, 0, &vhd->sul_central, sais_central_cb, 1); - } - - /* - * Recompute startable task platforms and broadcast to all sai-power, - * after there has been a change in tasks - */ - sais_platforms_with_tasks_pending(vhd); - - lwsl_notice("%s: exiting OK\n", __func__); - - return SAI_DB_RESULT_OK; -} - -sai_db_result_t -sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) -{ - char esc[96], cmd[256], event_uuid[33]; - sqlite3 *pdb = NULL; - lws_dll2_owner_t o; - struct lwsac *ac = NULL; - sai_task_t *task; - int ret; - - if (!task_uuid[0]) - return SAI_DB_RESULT_OK; - - lwsl_notice("%s: received request to rebuild last step of 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)) { - lwsl_err("%s: unable to open event-specific database\n", - __func__); - - return SAI_DB_RESULT_ERROR; - } - - lws_sql_purify(esc, task_uuid, sizeof(esc)); - lws_snprintf(cmd, sizeof(cmd), " and uuid='%s'", esc); - ret = lws_struct_sq3_deserialize(pdb, cmd, NULL, - lsm_schema_sq3_map_task, &o, &ac, 0, 1); - if (ret < 0 || !o.head) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - return SAI_DB_RESULT_ERROR; - } - - task = lws_container_of(o.head, sai_task_t, list); - - if (task->build_step > 0) { - lws_snprintf(cmd, sizeof(cmd), - "update tasks set build_step=%d where uuid='%s'", - task->build_step - 1, esc); - - ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - - lwsl_err("%s: %s: %s: fail\n", __func__, cmd, - sqlite3_errmsg(pdb)); - return SAI_DB_RESULT_ERROR; - } - } - - lwsac_free(&ac); - sais_event_db_close(vhd, &pdb); - - sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 0, 0); - - sais_task_stop_on_builders(vhd, task_uuid); - - lwsl_err("%s: scheduling sul_central to find a new task\n", __func__); - lws_sul_schedule(vhd->context, 0, &vhd->sul_central, sais_central_cb, 1); - - sais_platforms_with_tasks_pending(vhd); - - lwsl_notice("%s: exiting OK\n", __func__); - - return SAI_DB_RESULT_OK; -} - /* * Look for any task on any event that needs building on platform_name, if found * the caller must take responsibility to free pss->a.ac @@ -1122,8 +812,10 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, sai_task_t temp_task; int attempts = 0; - if (cb->busy) + if (cb->busy) { + lwsl_wsi_warn(pss->wsi, "::::::::::::: ABORTING task alloc due to BUSY on %s", cb->name); return 1; + } #if 0 if (cb->avail_slots <= 0) { @@ -1214,20 +906,24 @@ bail: return -1; } +#define MAX_BLOB 1024 + void sais_activity_cb(lws_sorted_usec_list_t *sul) { struct vhd *vhd = lws_container_of(sul, struct vhd, sul_activity); - char *p, *start, *end; - lws_usec_t now; - int cat, first = 1; struct lwsac *ac_events = NULL, *ac_tasks = NULL; lws_dll2_owner_t o_events, o_tasks; + char *p, *start, *end, *ast, s = 1; + int cat, first = 1; + lws_usec_t now; - p = start = malloc(8192); - if (!p) + ast = malloc(MAX_BLOB + LWS_PRE); + if (!ast) return; - end = start + 8192; + start = ast + LWS_PRE; + end = start + MAX_BLOB; + p = start; p += lws_snprintf(p, lws_ptr_diff_size_t(end, p), "{\"schema\":\"com.warmcat.sai.taskactivity\"," @@ -1239,6 +935,7 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) " and state != 3 and state != 4 and state != 5 and state != 7", NULL, lsm_schema_sq3_map_event, &o_events, &ac_events, 0, 100) >= 0 && o_events.head) { + lws_start_foreach_dll(struct lws_dll2 *, d, o_events.head) { sai_event_t *e = lws_container_of(d, sai_event_t, list); sqlite3 *pdb = NULL; @@ -1248,6 +945,7 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) " and (state = 1 or state = 2)", NULL, lsm_schema_sq3_map_task, &o_tasks, &ac_tasks, 0, 100) >= 0 && o_tasks.head) { + lws_start_foreach_dll(struct lws_dll2 *, dt, o_tasks.head) { sai_task_t *t = lws_container_of(dt, sai_task_t, list); @@ -1264,10 +962,16 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) if (!first) *p++ = ','; - p += lws_snprintf(p, lws_ptr_diff_size_t(end, p), - "{\"uuid\":\"%s\",\"cat\":%d}", - t->uuid, cat); + p += lws_snprintf(p, lws_ptr_diff_size_t(end, p), "{\"uuid\":\"%s\",\"cat\":%d}", t->uuid, cat); first = 0; + + if (lws_ptr_diff_size_t(end, p) < 100) { + /* we might start it, but it won't be the final frag here since we have JSON closure to do */ + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, start, lws_ptr_diff_size_t(p, start), + SAI_WEBSRV_PB__ACTIVITY, (s ? LWSSS_FLAG_SOM : 0)); + p = start; + s = 0; + } } lws_end_foreach_dll(dt); } lwsac_free(&ac_tasks); @@ -1280,14 +984,15 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) *p++ = ']'; *p++ = '}'; - if (!first) { - sais_websrv_broadcast(vhd->h_ss_websrv, start, - lws_ptr_diff_size_t(p, start)); - lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, - sais_activity_cb, 1 * LWS_US_PER_SEC); + if (!s) { /* ie, if we sent something, send the closing part of the JSON */ + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, start, + lws_ptr_diff_size_t(p, start), + SAI_WEBSRV_PB__ACTIVITY, LWSSS_FLAG_EOM); + + lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, sais_activity_cb, 1 * LWS_US_PER_SEC); } - free(start); + free(ast); } int @@ -1433,9 +1138,14 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for } if (!p) { /* no more steps */ + sai_uuid_list_t *u; + 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); + + if (sais_is_task_inflight(vhd, cb, temp_task->uuid, &u)) + sais_inflight_entry_destroy(u); ret = 0; goto bail; } diff --git a/src/server/s-webops.c b/src/server/s-webops.c new file mode 100644 index 0000000..05f3ac3 --- /dev/null +++ b/src/server/s-webops.c @@ -0,0 +1,354 @@ +/* + * 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 + * + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> +#include <assert.h> + +#include "s-private.h" + + +/* + * This is the only path to send things from server -> web + * + * 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. + * + * + * This is a bit tricky because the per sai-web buflist may be in the middle of + * a series of fragments for an existing message. We can't snipe our way in + * the middle and start dumping logs then. And, each sai-web connection may + * be in a different situation for ongoing existing messages. + * + * To solve this, we use lws_wsmsg_ apis to reassemble the various sources + * of messages using private buflists before emptying them into the upstream + * buflist. + */ + + +typedef struct { + const uint8_t *buf; + size_t len; + unsigned int ss_flags; + int reassembly_idx; +} sais_websrv_broadcast_t; + +static void +_sais_websrv_broadcast(struct lws_ss_handle *h, void *v) +{ + websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); + sais_websrv_broadcast_t *a = (sais_websrv_broadcast_t *)v; + unsigned int *pi = (unsigned int *)((const char *)a->buf - sizeof(int)); + + *pi = a->ss_flags; + + /* sai-web might not be taking it.. */ + + if (lws_buflist_total_len(&m->bl_srv_to_web) > (5u * 1024u * 1024u)) { + lwsl_ss_warn(h, "server->web buflist reached 5MB"); + lws_ss_start_timeout(h, 1); + return; + } + + if (lws_wsmsg_append(&m->bl_srv_to_web, + &m->private_heads[a->reassembly_idx], + a->buf - sizeof(int), + a->len + sizeof(int), a->ss_flags) < 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"); +} + +int +sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv, + const char *str, size_t len, + int reassembly_idx, unsigned int ss_flags) +{ + sais_websrv_broadcast_t a; + + a.buf = (const uint8_t *)str; /* LWS_PRE behind valid too */ + a.len = len; + a.ss_flags = ss_flags; + a.reassembly_idx = reassembly_idx; + + lws_ss_server_foreach_client(hsrv, _sais_websrv_broadcast, &a); + + return 0; +} + +struct sais_arg { + const char *uid; + int state; +}; + +static void +_sais_taskchange(struct lws_ss_handle *h, void *_arg) +{ + struct sais_arg *arg = (struct sais_arg *)_arg; + char tc[LWS_PRE + 128], *start = tc + LWS_PRE; + int n; + + n = lws_snprintf(start, sizeof(tc) - LWS_PRE, + "{\"schema\":\"sai-taskchange\", " + "\"event_hash\":\"%s\", \"state\":%d}", + arg->uid, arg->state); + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, start, (size_t)n, + SAI_WEBSRV_PB__GENERATED, + LWSSS_FLAG_SOM | LWSSS_FLAG_EOM) < 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) +{ + struct sais_arg *arg = (struct sais_arg *)_arg; + char tc[LWS_PRE + 128], *start = tc + LWS_PRE; + int n; + + n = lws_snprintf(start, sizeof(tc) - LWS_PRE, + "{\"schema\":\"sai-eventchange\", " + "\"event_hash\":\"%s\", \"state\":%d}", + arg->uid, arg->state); + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, start, (size_t)n, + SAI_WEBSRV_PB__GENERATED, + LWSSS_FLAG_SOM | LWSSS_FLAG_EOM) < 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 +sais_event_reset(struct vhd *vhd, const char *event_uuid) +{ + sqlite3 *pdb = NULL; + lws_dll2_owner_t o; + struct lwsac *ac = NULL; + char *err = NULL; + int ret; + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + return SAI_DB_RESULT_ERROR; + + if (lws_struct_sq3_deserialize(pdb, NULL, NULL, + lsm_schema_sq3_map_task, + &o, &ac, 0, 999) >= 0) { + + ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + return SAI_DB_RESULT_ERROR; + } + sqlite3_free(err); + + lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { + sai_task_t *t = lws_container_of(p, sai_task_t, list); + if (sais_task_clear_build_and_logs(vhd, t->uuid, 0) == SAI_DB_RESULT_BUSY) { + sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + return SAI_DB_RESULT_BUSY; + } + } lws_end_foreach_dll(p); + + ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + return SAI_DB_RESULT_ERROR; + } + sqlite3_free(err); + } + + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + + return SAI_DB_RESULT_OK; +} + +sai_db_result_t +sais_event_delete(struct vhd *vhd, const char *event_uuid) +{ + char qu[128], esc[96], pre[LWS_PRE + 128]; + struct lwsac *ac = NULL; + sqlite3 *pdb = NULL; + lws_dll2_owner_t o; + char *err = NULL; + size_t len; + int ret; + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb) == 0) { + if (lws_struct_sq3_deserialize(pdb, NULL, NULL, + lsm_schema_sq3_map_task, + &o, &ac, 0, 999) >= 0) { + + ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + return SAI_DB_RESULT_ERROR; + } + + lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { + sai_task_t *t = lws_container_of(p, sai_task_t, list); + + if (t->state != SAIES_WAITING && + t->state != SAIES_SUCCESS && + t->state != SAIES_FAIL && + t->state != SAIES_CANCELLED) + sais_task_cancel(vhd, t->uuid); + + } lws_end_foreach_dll(p); + + ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + return SAI_DB_RESULT_ERROR; + } + } + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + } + + lws_sql_purify(esc, event_uuid, sizeof(esc)); + lws_snprintf(qu, sizeof(qu), "delete from events where uuid='%s'", esc); + ret = sqlite3_exec(vhd->server.pdb, qu, NULL, NULL, &err); + if (ret != SQLITE_OK) { + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + lwsl_err("%s: evdel uuid %s, sq3 err %s\n", __func__, esc, err); + sqlite3_free(err); + return SAI_DB_RESULT_ERROR; + } + + sais_event_db_delete_database(vhd, event_uuid); + sais_eventchange(vhd->h_ss_websrv, event_uuid, SAIES_DELETED); + + len = (size_t)lws_snprintf(pre + LWS_PRE, sizeof(pre) - LWS_PRE, "{\"schema\":\"sai-overview\"}"); + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, pre + LWS_PRE, len, + SAI_WEBSRV_PB__GENERATED, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); + + return SAI_DB_RESULT_OK; +} + +sai_db_result_t +sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) +{ + sqlite3 *pdb = NULL; + lws_dll2_owner_t o; + struct lwsac *ac = NULL; + char *err = NULL; + int ret; + char filt[256], esc[96]; + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + return SAI_DB_RESULT_ERROR; + + lws_sql_purify(esc, platform, sizeof(esc)); + lws_snprintf(filt, sizeof(filt), " and platform='%s' and state=4", esc); + + if (lws_struct_sq3_deserialize(pdb, filt, NULL, + lsm_schema_sq3_map_task, + &o, &ac, 0, 999) >= 0) { + ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + return SAI_DB_RESULT_ERROR; + } + sqlite3_free(err); + + lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { + sai_task_t *t = lws_container_of(p, sai_task_t, list); + if (sais_task_clear_build_and_logs(vhd, t->uuid, 0) == SAI_DB_RESULT_BUSY) { + sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + return SAI_DB_RESULT_BUSY; + } + } lws_end_foreach_dll(p); + + ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); + if (ret != SQLITE_OK) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + if (ret == SQLITE_BUSY) + return SAI_DB_RESULT_BUSY; + return SAI_DB_RESULT_ERROR; + } + sqlite3_free(err); + } + + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + + return SAI_DB_RESULT_OK; +} + + + diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index c2c5d8a..d4b2760 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -1,7 +1,7 @@ /* * Sai server * - * 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 @@ -50,15 +50,6 @@ typedef struct sai_sul_retry_ctx { uint8_t op; /* SAIS_WS_WEBSRV_RX_... */ } sai_sul_retry_ctx_t; -typedef struct websrvss_srv { - struct lws_ss_handle *ss; - struct vhd *vhd; - /* ... application specific state ... */ - - struct lejp_ctx ctx; - struct lws_buflist *bltx; - unsigned int viewers; -} websrvss_srv_t; static lws_struct_map_t lsm_browser_taskreset[] = { LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "uuid"), @@ -166,67 +157,18 @@ reject: return 1; } -/* - * sais_webserv_broadcast allows us to queue to broadcast a message to all - * sai-web daemons that are connected to us. - * - * The queue is drained by websrvss_ws_tx() below. - * - * These messages are defined to all fit in a single fragment and will - * cause an assertion if they don't. - */ - -typedef struct { - const uint8_t *buf; - size_t len; -} sais_websrv_broadcast_t; - -static void -_sais_websrv_broadcast(struct lws_ss_handle *h, void *arg) -{ - websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); - sais_websrv_broadcast_t *a = (sais_websrv_broadcast_t *)arg; - - /* sai-web might not be taking it.. */ - - if (lws_buflist_total_len(&m->bltx) > 5000000u) { - lwsl_ss_warn(h, "server->web buflist reached 5MB"); - lws_ss_start_timeout(h, 1); - return; - } - - if (lws_buflist_append_segment(&m->bltx, a->buf, a->len) < 0) { - lwsl_err("%s: buflist append fail\n", __func__); - lws_ss_start_timeout(h, 1); - - return; - } - - if (lws_ss_request_tx(h)) - lwsl_ss_warn(h, "tx req fail"); -} - -void -sais_websrv_broadcast(struct lws_ss_handle *hsrv, const char *str, size_t len) -{ - sais_websrv_broadcast_t a; - - a.buf = (const uint8_t *)str; - a.len = len; - - lws_ss_server_foreach_client(hsrv, _sais_websrv_broadcast, &a); -} - int sais_list_builders(struct vhd *vhd) { - lws_dll2_owner_t db_builders_owner; - struct lwsac *ac = NULL; - char *p = vhd->json_builders, *end = p + sizeof(vhd->json_builders), + char json_builders[LWS_PRE + 1024], *start = json_builders + LWS_PRE, + *p = start, *end = p + sizeof(json_builders) - LWS_PRE, subsequent = 0; - lws_struct_serialize_t *js; + 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; size_t w; memset(&db_builders_owner, 0, sizeof(db_builders_owner)); @@ -244,6 +186,7 @@ sais_list_builders(struct vhd *vhd) "{\"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); @@ -257,20 +200,15 @@ sais_list_builders(struct vhd *vhd) 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; + 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; + builder_from_db->stay_on = live_builder->stay_on; } else - builder_from_db->online = 0; - - /* if (builder_from_db->power_managed) - lwsl_notice("%s: builder %s is power managed (stay: %d)\n", - __func__, builder_from_db->name, - builder_from_db->stay_on); */ + builder_from_db->online = 0; - builder_from_db->powering_up = 0; - builder_from_db->powering_down = 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); @@ -284,32 +222,52 @@ sais_list_builders(struct vhd *vhd) } } 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) { + 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; - if (lws_struct_json_serialize(js, (uint8_t *)p, - lws_ptr_diff_size_t(end, p), &w) != LSJS_RESULT_FINISH) { - lws_struct_json_serialize_destroy(&js); - goto bail; - } - p += w; + 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: + sais_websrv_broadcast_REQUIRES_LWS_PRE( + vhd->h_ss_websrv, start, + lws_ptr_diff_size_t(p, start), + SAI_WEBSRV_PB__GENERATED, + ss_flags); + 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), "]}"); + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, start, + lws_ptr_diff_size_t(p, start), + SAI_WEBSRV_PB__GENERATED, ss_flags); - // lwsl_notice("%s: Broadcasting builder list: %s\n", __func__, vhd->json_builders); - sais_websrv_broadcast(vhd->h_ss_websrv, vhd->json_builders, - lws_ptr_diff_size_t(p, vhd->json_builders)); - + // lwsl_notice("%s: Broadcasting builder list: %s\n", __func__, start); lwsac_free(&ac); return 0; @@ -318,69 +276,7 @@ bail: return 1; } -struct sais_arg { - const char *uid; - int state; -}; - -static void -_sais_taskchange(struct lws_ss_handle *h, void *_arg) -{ - websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); - struct sais_arg *arg = (struct sais_arg *)_arg; - char tc[128]; - int n; - - n = lws_snprintf(tc, sizeof(tc), "{\"schema\":\"sai-taskchange\", " - "\"event_hash\":\"%s\", \"state\":%d}", - arg->uid, arg->state); - - if (lws_buflist_append_segment(&m->bltx, (uint8_t *)tc, (unsigned int)n) < 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) -{ - websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); - struct sais_arg *arg = (struct sais_arg *)_arg; - char tc[128]; - int n; - - n = lws_snprintf(tc, sizeof(tc), "{\"schema\":\"sai-eventchange\", " - "\"event_hash\":\"%s\", \"state\":%d}", - arg->uid, arg->state); - - if (lws_buflist_append_segment(&m->bltx, (uint8_t *)tc, (unsigned int)n) < 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); -} static void sum_viewers_cb(struct lws_ss_handle *h, void *arg) @@ -389,181 +285,7 @@ sum_viewers_cb(struct lws_ss_handle *h, void *arg) *(unsigned int *)arg += m_client->viewers; } -static sai_db_result_t -sais_event_reset(struct vhd *vhd, const char *event_uuid) -{ - sqlite3 *pdb = NULL; - lws_dll2_owner_t o; - struct lwsac *ac = NULL; - char *err = NULL; - int ret; - - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) - return SAI_DB_RESULT_ERROR; - - if (lws_struct_sq3_deserialize(pdb, NULL, NULL, - lsm_schema_sq3_map_task, - &o, &ac, 0, 999) >= 0) { - - ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - return SAI_DB_RESULT_ERROR; - } - sqlite3_free(err); - - lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { - sai_task_t *t = lws_container_of(p, sai_task_t, list); - if (sais_task_clear_build_and_logs(vhd, t->uuid, 0) == SAI_DB_RESULT_BUSY) { - sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - return SAI_DB_RESULT_BUSY; - } - } lws_end_foreach_dll(p); - - ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - return SAI_DB_RESULT_ERROR; - } - sqlite3_free(err); - } - - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - - return SAI_DB_RESULT_OK; -} - -sai_db_result_t -sais_event_delete(struct vhd *vhd, const char *event_uuid) -{ - sqlite3 *pdb = NULL; - lws_dll2_owner_t o; - struct lwsac *ac = NULL; - char *err = NULL; - int ret; - char qu[128], esc[96]; - - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb) == 0) { - if (lws_struct_sq3_deserialize(pdb, NULL, NULL, - lsm_schema_sq3_map_task, - &o, &ac, 0, 999) >= 0) { - - ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - return SAI_DB_RESULT_ERROR; - } - - lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { - sai_task_t *t = lws_container_of(p, sai_task_t, list); - - if (t->state != SAIES_WAITING && - t->state != SAIES_SUCCESS && - t->state != SAIES_FAIL && - t->state != SAIES_CANCELLED) - sais_task_cancel(vhd, t->uuid); - - } lws_end_foreach_dll(p); - - ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - return SAI_DB_RESULT_ERROR; - } - } - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - } - - lws_sql_purify(esc, event_uuid, sizeof(esc)); - lws_snprintf(qu, sizeof(qu), "delete from events where uuid='%s'", esc); - ret = sqlite3_exec(vhd->server.pdb, qu, NULL, NULL, &err); - if (ret != SQLITE_OK) { - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - lwsl_err("%s: evdel uuid %s, sq3 err %s\n", __func__, esc, err); - sqlite3_free(err); - return SAI_DB_RESULT_ERROR; - } - - sais_event_db_delete_database(vhd, event_uuid); - sais_eventchange(vhd->h_ss_websrv, event_uuid, SAIES_DELETED); - sais_websrv_broadcast(vhd->h_ss_websrv, - "{\"schema\":\"sai-overview\"}", 25); - - return SAI_DB_RESULT_OK; -} - -static sai_db_result_t -sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) -{ - sqlite3 *pdb = NULL; - lws_dll2_owner_t o; - struct lwsac *ac = NULL; - char *err = NULL; - int ret; - char filt[256], esc[96]; - - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) - return SAI_DB_RESULT_ERROR; - - lws_sql_purify(esc, platform, sizeof(esc)); - lws_snprintf(filt, sizeof(filt), " and platform='%s' and state=4", esc); - - if (lws_struct_sq3_deserialize(pdb, filt, NULL, - lsm_schema_sq3_map_task, - &o, &ac, 0, 999) >= 0) { - ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - return SAI_DB_RESULT_ERROR; - } - sqlite3_free(err); - - lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { - sai_task_t *t = lws_container_of(p, sai_task_t, list); - if (sais_task_clear_build_and_logs(vhd, t->uuid, 0) == SAI_DB_RESULT_BUSY) { - sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - return SAI_DB_RESULT_BUSY; - } - } lws_end_foreach_dll(p); - - ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); - if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - if (ret == SQLITE_BUSY) - return SAI_DB_RESULT_BUSY; - return SAI_DB_RESULT_ERROR; - } - sqlite3_free(err); - } - - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - return SAI_DB_RESULT_OK; -} static lws_ss_state_return_t @@ -788,29 +510,57 @@ 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; - size_t fsl = lws_buflist_next_segment_len(&m->bltx, NULL); - char som, eom; - int used; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->bl_srv_to_web), depi; + char som, som1, eom, final = 1; + size_t fsl, used; - if (!m->bltx) + if (!m->bl_srv_to_web) return LWSSSSRET_TX_DONT_SEND; - used = lws_buflist_fragment_use(&m->bltx, buf, *len, &som, &eom); + 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 ((size_t)used < fsl) - eom = 0; /* because we still be back */ + if (used < fsl || !(depi & LWSSS_FLAG_EOM)) /* we saved SS flags at the start of the buf */ + final = 0; + + *len = used; + *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - *flags = (som ? LWSSS_FLAG_SOM : 0) | (eom ? LWSSS_FLAG_EOM : 0); - *len = (size_t)used; + // lwsl_ss_notice(m->ss, "Sending %d srv->web: som %d, som1 %d, depi %d, ssflags %d", (int)*len, som, som1, depi, (int)*flags); + // lwsl_hexdump_notice(buf, *len); - if (m->bltx) + 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) @@ -824,7 +574,9 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, case LWSSSCS_DISCONNECTED: { unsigned int total_viewers = 0; - lws_buflist_destroy_all_segments(&m->bltx); + 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 */ @@ -843,6 +595,7 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, 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); @@ -850,6 +603,7 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, } } lws_end_foreach_dll(p); } + break; } case LWSSSCS_CREATING: @@ -857,7 +611,6 @@ websrvss_srv_state(void *userobj, void *sh, lws_ss_constate_t state, return lws_ss_request_tx(m->ss); case LWSSSCS_CONNECTED: - // lwsl_warn("%s: resending builders because CONNECTED\n", __func__); sais_list_builders(m->vhd); break; case LWSSSCS_ALL_RETRIES_FAILED: diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 6b3dcfe..16c2507 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -1,7 +1,7 @@ /* * Sai server - ./src/server/s-ws-builder.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 @@ -25,6 +25,7 @@ #include <libwebsockets.h> #include <string.h> #include <signal.h> +#include <assert.h> #include <time.h> #include "s-private.h" @@ -85,7 +86,7 @@ sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) { struct vhd *vhd = lws_container_of(sul, struct vhd, sul_logcache); sais_logcache_pertask_t *lcpt; - char event_uuid[33], sw[192]; + char event_uuid[33], sw[192 + LWS_PRE]; sqlite3 *pdb = NULL; sai_log_t *hlog; char *err; @@ -101,6 +102,7 @@ sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) sai_task_uuid_to_event_uuid(event_uuid, lcpt->uuid); + pdb = NULL; if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { /* @@ -141,9 +143,13 @@ sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) * something changed (event_hash is actually the task hash) */ - n = lws_snprintf(sw, sizeof(sw), "{\"schema\":\"sai-tasklogs\"," + n = lws_snprintf(sw + LWS_PRE, sizeof(sw) - LWS_PRE, + "{\"schema\":\"sai-tasklogs\"," "\"event_hash\":\"%s\"}", lcpt->uuid); - sais_websrv_broadcast(vhd->h_ss_websrv, sw, (unsigned int)n); + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, sw + LWS_PRE, + (unsigned int)n, + SAI_WEBSRV_PB__LOGS, + LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); /* * Destroy the whole task-specific cache, it will regenerate @@ -393,8 +399,16 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) lwsac_free(&ac); } - lws_snprintf(q, sizeof(q), "UPDATE builders SET online=0 WHERE name='%s'", cb->name); - sai_sqlite3_statement(vhd->server.pdb, q, "set builder offline"); + /* drop any inflight task information for this builder */ + + lws_start_foreach_dll_safe(struct lws_dll2 *, pif, pif1, + cb->inflight_owner.head) { + sai_uuid_list_t *ul = lws_container_of(pif, sai_uuid_list_t, list); + + sais_inflight_entry_destroy(ul); + + } lws_end_foreach_dll_safe(pif, pif1); + const char *dot = strchr(cb->name, '.'); if (dot) { @@ -412,6 +426,8 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) lws_dll2_remove(&cb->sai_plat_list); free(cb); + + // assert(0); } } lws_end_foreach_dll_safe(p, p1); } @@ -428,10 +444,12 @@ sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name) /* * Server received a communication from a builder + * + * buf is lws callback `in` which has LWS_PRE already set aside */ int -sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) +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; @@ -500,6 +518,9 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b // lwsl_hexdump_notice(buf, bl); if (m == LEJP_CONTINUE) { + if (pss->a.top_schema_index == SAIM_WSSCH_BUILDER_LOADREPORT) + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, (const char *)buf, bl, + SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, ss_flags); pss->frag = 1; return 0; } @@ -552,7 +573,7 @@ handle: if (live_cb) { /* Already exists (reconnect), just update dynamic info */ lwsl_err("%s: found live builder for %s\n", __func__, build->name); - live_cb->wsi = pss->wsi; + live_cb->wsi = pss->wsi; lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip)); lws_strncpy(live_cb->sai_hash, build->sai_hash, sizeof(live_cb->sai_hash)); @@ -578,21 +599,21 @@ handle: char *p_str = (char *)(live_cb + 1); memset(live_cb, 0, sizeof(*live_cb)); - live_cb->name = p_str; + live_cb->name = p_str; memcpy(p_str, build->name, nlen); - live_cb->platform = p_str + nlen; + live_cb->platform = p_str + nlen; memcpy(p_str + nlen, build->platform, plen); lws_strncpy(live_cb->sai_hash, build->sai_hash, sizeof(live_cb->sai_hash)); lws_strncpy(live_cb->lws_hash, build->lws_hash, sizeof(live_cb->lws_hash)); - live_cb->windows = build->windows; - live_cb->avail_slots = 1; /* default */ - live_cb->avail_mem_kib = (unsigned int)-1; - live_cb->avail_sto_kib = (unsigned int)-1; - live_cb->s_avail_slots = live_cb->avail_slots; - live_cb->wsi = pss->wsi; - live_cb->online = 1; + live_cb->windows = build->windows; + live_cb->avail_slots = 1; /* default */ + live_cb->avail_mem_kib = (unsigned int)-1; + live_cb->avail_sto_kib = (unsigned int)-1; + live_cb->s_avail_slots = live_cb->avail_slots; + live_cb->wsi = pss->wsi; + live_cb->online = 1; lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip)); lws_dll2_add_tail(&live_cb->sai_plat_list, &vhd->server.builder_owner); } @@ -666,14 +687,15 @@ bail: if (log->finished) { sai_plat_t *cb; - char builder_name[128], esc_uuid[129], q[128], - event_uuid[33]; + // 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 and update our * tracking of its state */ + 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'; @@ -685,17 +707,18 @@ bail: NULL) == SQLITE_OK && builder_name[0]) { cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); if (cb) { - sai_uuid_list_t *sul; + // 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->busy = 0; + __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->busy = 0; +#if 0 lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, cb->inflight_owner.head) { sul = lws_container_of(d, sai_uuid_list_t, list); if (!strcmp(sul->uuid, log->task_uuid)) { @@ -703,6 +726,7 @@ bail: break; } } lws_end_foreach_dll_safe(d, d1); +#endif cb->s_avail_slots = cb->avail_slots; cb->s_inflight_count = (int)cb->inflight_owner.count; @@ -713,31 +737,41 @@ bail: } sais_event_db_close(vhd, &pdb); } + /* - * We have reached the end of the logs for this task + * We have reached the end of the logs for this task step */ sais_dump_logs_to_db(&vhd->sul_logcache); +#if 0 + /* + * Remove us from the inflight list + */ + + if (sais_is_task_inflight(vhd, cb, log->task_uuid, &u)) + sais_inflight_entry_destroy(u); +#endif + 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\n", __func__); + lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_SUCCESS: %s\n", __func__, log->task_uuid); } else { n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL\n", __func__); + lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL: %s\n", __func__, log->task_uuid); } } else if (log->finished & 0x2000) { n = SAIES_CANCELLED; - lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED\n", __func__); + lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED: %s\n", __func__, log->task_uuid); } else { n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_FAIL\n", __func__); + 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, @@ -781,25 +815,47 @@ bail: switch (rej->reason) { case SAI_TASK_REASON_ACCEPTED: - lwsl_notice("%s: SAI_TASK_REASON_ACCEPTED\n", __func__); - pss->first_log_timestamp = lws_now_secs(); + lwsl_notice("%s: SAI_TASK_REASON_ACCEPTED: %s\n", __func__, rej->task_uuid); + { + char event_uuid[33]; + sqlite3 *pdb = NULL; + int build_step = -1; + + sai_task_uuid_to_event_uuid(event_uuid, rej->task_uuid); + if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + char q[128], esc_uuid[129]; + + lws_sql_purify(esc_uuid, rej->task_uuid, sizeof(esc_uuid)); + 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; + 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, NULL, NULL, rej->task_uuid, SAIES_BEING_BUILT, 0, 0)) break; /* leave the uuid listed until step completed */ break; case SAI_TASK_REASON_DUPE: - lwsl_notice("%s: SAI_TASK_REASON_DUPE\n", __func__); - // do_remove_uuid = 1; + 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\n", __func__); + lwsl_notice("%s: SAI_TASK_REASON_BUSY: Set busy: %s\n", __func__, rej->task_uuid); do_remove_uuid = 1; cb->busy = 1; break; case SAI_TASK_REASON_DESTROYED: - lwsl_notice("%s: SAI_TASK_REASON_DESTROYED\n", __func__); + lwsl_notice("%s: SAI_TASK_REASON_DESTROYED: Clear busy: %s\n", __func__, rej->task_uuid); do_remove_uuid = 1; + cb->busy = 0; break; } @@ -816,8 +872,8 @@ bail: // sais_task_clear_build_and_logs(vhd, rej->task_uuid, 31); - cb->s_avail_slots = cb->avail_slots; - cb->s_inflight_count = (int)cb->inflight_owner.count; + 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)); @@ -839,7 +895,8 @@ bail: // lwsl_notice("%s: write failed\n", __func__); // } // lwsl_wsi_user(pss->wsi, "SAIM_WSSCH_BUILDER_LOADREPORT broadcasting\n"); - sais_websrv_broadcast(vhd->h_ss_websrv, (const char *)buf, bl); + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, (const char *)buf, bl, + SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, ss_flags); break; case SAIM_WSSCH_BUILDER_ARTIFACT: @@ -1146,9 +1203,18 @@ bail: LWS_ARRAY_SIZE(lsm_schema_build_metric), 0, (void *)metric); if (js) { - int n = lws_struct_json_serialize(js, buf, sizeof(buf), &used); - if (n >= 0) - sais_websrv_broadcast(vhd->h_ss_websrv, (const char *)buf, used); + 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: + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, (const char *)buf, used, + SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); + break; + } lws_struct_json_serialize_destroy(&js); } } diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c index 5839dd4..2ec0151 100644 --- a/src/web/w-websrv.c +++ b/src/web/w-websrv.c @@ -353,4 +353,43 @@ const lws_ss_info_t ssi_saiw_websrv = { .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 35e6c25..d2ac344 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -302,14 +302,14 @@ 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: n %d, o.head %p\n", __func__, n, o.head); + 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_info("%s: browser ws asked for task hash: %s, plat %s\n", + 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 */ @@ -348,15 +348,19 @@ 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); - if (n < 0 || !o.head) { - // lwsl_notice("%s: no result\n", __func__); - goto bail; - } + lwsl_notice("%s: WWWWWWWWWWW -- actual event n %d, o.head %p\n", __func__, n, o.head); - sch->logsub = !!logsub; - sch->one_event = lws_container_of(o.head, sai_event_t, list); + if (n < 0 || !o.head) + /* + * It's OK if the parent event is not visible in the current + * filtered view, we can still update the task state where it + * appears inside other visible events + */ + sch->one_event = NULL; + else + sch->one_event = lws_container_of(o.head, sai_event_t, list); - // lwsl_warn("%s: doing WSS_PREPARE_BUILDER_SUMMARY\n", __func__); + sch->logsub = !!logsub; saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); @@ -680,10 +684,11 @@ again: // lwsl_notice("%s: send_state %d, pss %p, wsi %p\n", __func__, // pss->send_state, pss, pss->wsi); - if (pss->sched.count) + // lwsl_warn("%s: pss->sched.count %d, pss->sched.head %p\n", __func__, pss->sched.count, pss->sched.head); + + sch = NULL; + if (pss->sched.head) sch = lws_container_of(pss->sched.head, saiw_scheduled_t, list); - else - sch = NULL; switch (pss->send_state) { case WSS_IDLE1: @@ -694,7 +699,7 @@ again: * If so, let's prioritize that first... */ - if ((!pss->sched.count || !pss->toggle_favour_sch) && + if ((!pss->sched.head || !pss->toggle_favour_sch) && pss->subs_list.owner) { sch = NULL; @@ -1226,16 +1231,16 @@ b_finish: * when we go out of scope... */ - lwsl_info("%s: PREPARE_TASKINFO: one_task %p\n", __func__, sch->one_task); + lwsl_warn("%s: wwwwwwwwwwww PREPARE_TASKINFO: one_task %p\n", __func__, sch->one_task); - task_reply.event = sch->one_event; - task_reply.task = sch->one_task; - sch->one_task->rebuildable = (sch->one_task->state == SAIES_FAIL || - sch->one_task->state == SAIES_CANCELLED) && - (lws_now_secs() - (sch->one_task->started + + task_reply.event = sch->one_event; + task_reply.task = sch->one_task; + sch->one_task->rebuildable = (sch->one_task->state == SAIES_FAIL || + sch->one_task->state == SAIES_CANCELLED) && + (lws_now_secs() - (sch->one_task->started + (sch->one_task->duration / 1000000)) < 24 * 3600); - task_reply.auth_secs = (int)(pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0); - task_reply.authorized = pss->authorized; + task_reply.auth_secs = (int)(pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0); + task_reply.authorized = pss->authorized; lws_strncpy(task_reply.auth_user, pss->auth_user, sizeof(task_reply.auth_user)); @@ -1307,6 +1312,12 @@ b_finish: lwsl_notice("%s: taskinfo: empty json\n", __func__); return 0; } + + 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"); + break; case WSS_SEND_ARTIFACT_INFO: @@ -1412,43 +1423,4 @@ saiw_browser_state_changed(struct pss *pss, int established) saiw_update_viewer_count(pss->vhd); } -/* - * 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: 18ms (vhost etag hits: 0%, cache hits: 0%)