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 / scripts / etc-rc.d-sai_builder-OpenBSD
Author[]Andy Green <andy@warmcat.com> 2025-11-12 08:23 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-11-15 12:18 UTC
Treea4bfaeb765e1d778e1975dd151cd909cee67343b   Raw Patch
 
queue: convert sch in web
queue: convert sch in web
diff --git a/src/builder/b-power.c b/src/builder/b-power.c index 6d0f9f3..d82be98 100644 --- a/src/builder/b-power.c +++ b/src/builder/b-power.c @@ -63,6 +63,7 @@ saib_reassess_idle_situation() struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list); + if (xns->task) lwsl_notice("%s: ongoing task: %s\n", __func__, xns->task->uuid); @@ -258,6 +259,29 @@ sul_do_suspend_cb(lws_sorted_usec_list_t *sul) #endif } +void +sul_shutdown_cb(lws_sorted_usec_list_t *sul) +{ + int fd = saib_suspender_get_pipe(); + uint8_t te = 0; + ssize_t n; + + lwsl_warn("%s: device shutting down\n", __func__); + + n = write(fd, &te, 1); + + if (n != 1) + lwsl_err("%s: shutdown request failed\n", __func__); + +#if defined(WIN32) + Sleep(40000); +#else + sleep(40); +#endif + + lwsl_err("%s: shutdown didn't happen\n", __func__); +} + /* * The grace time is up, ask for the suspend */ @@ -320,6 +344,14 @@ sul_idle_cb(lws_sorted_usec_list_t *sul) if (!builder.url_sai_power) return; + /* + * We're planning to get ourselves turned off after we have shutdown + * cleanly. + * + * Send the request to sai-power to turn us off after 35s and then + * request our suspender process to shutdown the device. + */ + snprintf(path, sizeof(path) - 1, "%s/auto-power-off/%s", builder.url_sai_power, builder.host); @@ -333,6 +365,13 @@ sul_idle_cb(lws_sorted_usec_list_t *sul) if (lws_ss_request_tx(builder.ss_power_off)) lwsl_ss_warn(builder.ss_power_off, "Unable to request tx"); + + /* allow time for the sai-power transaction to happen */ + + lws_sul_schedule(builder.context, 0, &builder.sul_do_shutdown, + sul_shutdown_cb, 2 * LWS_US_PER_SEC); + + /* let event loop continue for a couple of seconds, then shutdown */ } int diff --git a/src/builder/b-private.h b/src/builder/b-private.h index d868c6c..216f9c8 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -81,19 +81,6 @@ struct saib_opaque_spawn { #define SAI_CLEANUP_JOBS_INTERVAL_US (60 * 60 * LWS_US_PER_SEC) #define SAI_CLEANUP_JOB_DIR_MIN_AGE_SECS (24ull * 3600u) -typedef enum { - PHASE_IDLE, - - PFL_FIRST = 128, - - PHASE_START_ATTACH = PFL_FIRST | 1, - PHASE_SUMM_PLATFORMS = 2, - - PHASE_BUILDING - -} cursor_phase_t; - - struct saib_ws_pss; @@ -106,8 +93,6 @@ enum nsstate { NSSTATE_FAILED, }; - - /* * This represents this builder process as a whole */ @@ -130,6 +115,7 @@ struct sai_builder { lws_sorted_usec_list_t sul_idle; lws_sorted_usec_list_t sul_do_suspend; + lws_sorted_usec_list_t sul_do_shutdown; lws_sorted_usec_list_t sul_stay; lws_sorted_usec_list_t sul_cleanup_jobs; diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index f48e4be..5bf6910 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -595,9 +595,17 @@ int main(int argc, const char **argv) } saib_power_init(); + #if defined(__linux__) - if (saib_suspender_fork(argv[0])) - return 1; + if (builder.power_off_type && + !strcmp(builder.power_off_type, "suspend") && + saib_suspender_fork(argv[0])) + return 1; +#endif + +#if defined(__NetBSD__) + if (saib_suspender_fork(argv[0])) + return 1; #endif while (!lws_service(builder.context, 0) && !interrupted) diff --git a/src/builder/b-suspender.c b/src/builder/b-suspender.c index 4f7b60e..7dd910c 100644 --- a/src/builder/b-suspender.c +++ b/src/builder/b-suspender.c @@ -249,6 +249,9 @@ saib_suspender_start(void) n = read(0, &d, 1); lwsl_notice("%s: suspend process read returned %d\n", __func__, (int)n); +#if defined(__APPLE__) + sleep(1); +#endif if (n <= 0) continue; diff --git a/src/common/c-sqlite3.c b/src/common/c-sqlite3.c new file mode 100644 index 0000000..bf3d4a1 --- /dev/null +++ b/src/common/c-sqlite3.c @@ -0,0 +1,239 @@ +/* + * Sai common utils + * + * Copyright (C) 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 <assert.h> + +#include "include/private.h" + +int +sai_event_db_ensure_open(struct lws_context *cx, lws_dll2_owner_t *sqlite3_cache, + const char *sqlite3_path_lhs, 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, 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", + sqlite3_path_lhs, saf); + + if (lws_struct_sq3_open(cx, 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, sqlite3_cache); + + return 0; +} + + +void +sai_event_db_close(lws_dll2_owner_t *sqlite3_cache, sqlite3 **ppdb) +{ + sais_sqlite_cache_t *sc; + + if (!*ppdb) + return; + + /* look for him in the cache */ + + lws_start_foreach_dll(struct lws_dll2 *, p, 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 +sai_event_db_close_all_now(lws_dll2_owner_t *sqlite3_cache) +{ + sais_sqlite_cache_t *sc; + + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + 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 +sai_event_db_delete_database(const char *sqlite3_path_lhs, 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", + 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", + 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", + 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; +} + +/* 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; +} diff --git a/src/common/c-utils.c b/src/common/c-utils.c index f3d85f9..96d0f40 100644 --- a/src/common/c-utils.c +++ b/src/common/c-utils.c @@ -143,8 +143,8 @@ sai_ss_serialize_queue_helper(struct lws_ss_handle *h, sizeof(buf) - LWS_PRE, &w); sai_ss_queue_frag_on_buflist_REQUIRES_LWS_PRE(h, buflist, - buf + LWS_PRE, w, (fi ? LWSSS_FLAG_SOM : 0) | - (r == LSJS_RESULT_FINISH ? LWSSS_FLAG_EOM : 0)); + buf + LWS_PRE, w, (unsigned int)((fi ? LWSSS_FLAG_SOM : 0) | + (r == LSJS_RESULT_FINISH ? LWSSS_FLAG_EOM : 0))); fi = 0; } while (r == LSJS_RESULT_CONTINUE); @@ -208,216 +208,3 @@ sai_ss_tx_from_buflist_helper(struct lws_ss_handle *ss, struct lws_buflist **buf return LWSSSSRET_OK; } - -int -sai_event_db_ensure_open(struct lws_context *cx, lws_dll2_owner_t *sqlite3_cache, - const char *sqlite3_path_lhs, 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, 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", - sqlite3_path_lhs, saf); - - if (lws_struct_sq3_open(cx, 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, sqlite3_cache); - - return 0; -} - - -void -sai_event_db_close(lws_dll2_owner_t *sqlite3_cache, sqlite3 **ppdb) -{ - sais_sqlite_cache_t *sc; - - if (!*ppdb) - return; - - /* look for him in the cache */ - - lws_start_foreach_dll(struct lws_dll2 *, p, 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 -sai_event_db_close_all_now(lws_dll2_owner_t *sqlite3_cache) -{ - sais_sqlite_cache_t *sc; - - lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, - 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 -sai_event_db_delete_database(const char *sqlite3_path_lhs, 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", - 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", - 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", - 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; -} - -/* 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; -} diff --git a/src/power/p-http-api.c b/src/power/p-http-api.c index 8e18fef..fa14f8b 100644 --- a/src/power/p-http-api.c +++ b/src/power/p-http-api.c @@ -355,18 +355,23 @@ power_off: */ needs[0] = '\0'; - lws_start_foreach_dll(struct lws_dll2 *, px1, sp->dependencies_owner.head) { - saip_server_plat_t *sp1 = lws_container_of(px1, saip_server_plat_t, dependencies_list); + lws_start_foreach_dll(struct lws_dll2 *, px1, + sp->dependencies_owner.head) { + saip_server_plat_t *sp1 = lws_container_of(px1, + saip_server_plat_t, dependencies_list); if (sp1->needed) - lws_snprintf(needs, sizeof(needs) - 1 - strlen(needs), "%s ", sp1->name); + lws_snprintf(needs, + sizeof(needs) - 1 - strlen(needs), + "%s ", sp1->name); } lws_end_foreach_dll(px1); if (needs[0] || sp->needed) { - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "NAK: %s needed: %d, deps needed: '%s'", - pn, sp->needed, needs); + g->size = (size_t)lws_snprintf(g->payload, + sizeof(g->payload), + "NAK: %s needed: %d, deps needed: '%s'", + pn, sp->needed, needs); goto bail; } } @@ -377,13 +382,16 @@ power_off: lws_sul_schedule(lws_ss_cx_from_user(g), 0, &sp->sul_delay_off, saip_sul_action_power_off, - 3 * LWS_USEC_PER_SEC); + SAI_POWERDOWN_HOLDOFF_US); - lwsl_warn("%s: scheduled powering off host %s\n", - __func__, sp->host); + lwsl_warn("%s: scheduled powering off host %s in %ds\n", + __func__, sp->host, + (int)(SAI_POWERDOWN_HOLDOFF_US / LWS_USEC_PER_SEC)); g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "ACK: Scheduled powering off host %s", sp->host); + "ACK: Scheduled powering off host %s in %ds", + sp->host, + (int)(SAI_POWERDOWN_HOLDOFF_US / LWS_USEC_PER_SEC)); sp->stay = 0; /* reset any manual power up */ } diff --git a/src/power/p-private.h b/src/power/p-private.h index 9911d15..332f94e 100644 --- a/src/power/p-private.h +++ b/src/power/p-private.h @@ -38,19 +38,7 @@ #endif #include <pthread.h> -#define SAI_IDLE_GRACE_US (20 * LWS_US_PER_SEC) - -typedef enum { - PHASE_IDLE, - - PFL_FIRST = 128, - - PHASE_START_ATTACH = PFL_FIRST | 1, - PHASE_SUMM_PLATFORMS = 2, - - PHASE_BUILDING - -} cursor_phase_t; +#define SAI_POWERDOWN_HOLDOFF_US (50 * LWS_US_PER_SEC) typedef struct tasmota_data { unsigned int voltage_v; diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index 576dfa0..8595ed7 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -17,6 +17,7 @@ set(SRCS s-webops.c s-resource.c ../common/c-utils.c + ../common/c-sqlite3.c ../common/struct-metadata.c ) diff --git a/src/web/CMakeLists.txt b/src/web/CMakeLists.txt index cd167a5..e91bf34 100644 --- a/src/web/CMakeLists.txt +++ b/src/web/CMakeLists.txt @@ -10,6 +10,7 @@ set(SRCS w-ws-server.c w-ws-browser.c ../common/c-utils.c + ../common/c-sqlite3.c ../common/struct-metadata.c ) diff --git a/src/web/w-comms.c b/src/web/w-comms.c index d113ef8..13d5efb 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -18,8 +18,7 @@ * 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). + * This ws interface is provides the 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 @@ -96,16 +95,6 @@ saiw_task_cancel(struct vhd *vhd, const char *task_uuid) } int -saiw_sched_destroy(struct lws_dll2 *d, void *user) -{ - saiw_scheduled_t *sch = lws_container_of(d, saiw_scheduled_t, list); - - saiw_dealloc_sched(sch); - - return 0; -} - -int sai_get_head_status(struct vhd *vhd, const char *projname) { struct lwsac *ac = NULL; @@ -587,7 +576,8 @@ http_resp: * It means, logout then */ - n = lws_snprintf(temp, sizeof(temp), "__Host-sai_jwt=deleted;" + n = lws_snprintf(temp, sizeof(temp), + "__Host-sai_jwt=deleted;" "HttpOnly;" "Secure;" "SameSite=strict;" @@ -596,8 +586,9 @@ http_resp: sr = "x/.."; - if (lws_add_http_header_by_token(wsi, WSI_TOKEN_HTTP_SET_COOKIE, - (uint8_t *)temp, n, &p, end)) { + if (lws_add_http_header_by_token(wsi, + WSI_TOKEN_HTTP_SET_COOKIE, + (uint8_t *)temp, n, &p, end)) { lwsl_err("%s: failed to add JWT cookie header\n", __func__); return 1; } @@ -697,7 +688,8 @@ back: return 0; final: - lwsl_notice("%s: auth failed, login_form %d\n", __func__, pss->login_form); + lwsl_notice("%s: auth failed, login_form %d\n", + __func__, pss->login_form); /* * Auth failed, go back to / */ @@ -745,7 +737,7 @@ clean_spa: /* * This protocol is for browsers on /browse... URLs. * Builders connect on /builder... URLs and should be handled - * by a different protocol. Explicitly reject them here. + * by sai-server. Explicitly reject them here. * * Returning 0 accepts the connection for this protocol. * Returning non-zero rejects it. @@ -851,6 +843,7 @@ clean_spa: if (!strncmp(tbuf, "task=", 5)) { lws_strncpy(pss->specific_task, tbuf + 5, sizeof(pss->specific_task)); pss->specificity = SAIM_SPECIFIC_TASK; + saiw_broadcast_logs_batch(vhd, pss); } if (!strncmp(tbuf, "h=", 2)) { memcpy(pss->specific_ref, "refs/heads/", 11); @@ -897,8 +890,8 @@ clean_spa: lws_buflist_destroy_all_segments(&pss->raw_tx); saiw_browser_state_changed(pss, 0); lws_dll2_remove(&pss->subs_list); + lws_sul_cancel(&pss->sul_logcache); - lws_dll2_foreach_safe(&pss->sched, NULL, saiw_sched_destroy); lwsac_free(&pss->logs_ac); break; @@ -921,12 +914,34 @@ clean_spa: break; case LWS_CALLBACK_SERVER_WRITEABLE: - if (!vhd) { - lwsl_notice("%s: no vhd\n", __func__); + if (!vhd || !pss->raw_tx) break; + + { + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&pss->raw_tx), depi = *pi; + char som, eom, rb[1200]; + int used, final = 1; + size_t fsl = lws_buflist_next_segment_len(&pss->raw_tx, NULL); + + /* this is the only buflist user on pss->raw_tx */ + used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, sizeof(rb), &som, &eom); + if (!used) + return 0; + if (used < (int)fsl || (depi & LWS_WRITE_NO_FIN)) + final = 0; + + if (lws_write(pss->wsi, (uint8_t *)rb + ((size_t)som * sizeof(int)), + (size_t)used - ((size_t)som * sizeof(int)), + (lws_ws_sending_multifragment(pss->wsi) ? LWS_WRITE_CONTINUATION : LWS_WRITE_TEXT) | + (!final * LWS_WRITE_NO_FIN)) < 0) { + lwsl_wsi_err(pss->wsi, "attempt to write %d failed", (int)used - (int)sizeof(int)); + + return -1; + } } - saiw_ws_json_tx_browser(vhd, pss, buf, sizeof(buf)); + if (pss->raw_tx) + lws_callback_on_writable(pss->wsi); break; default: diff --git a/src/web/w-private.h b/src/web/w-private.h index eff812c..7717cd8 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -46,22 +46,6 @@ typedef struct sai_platform { /* build and name over-allocated here */ } sai_platform_t; -typedef enum { - WSS_IDLE1, - WSS_IDLE2, - WSS_IDLE3, - WSS_PREPARE_OVERVIEW, - WSS_SEND_OVERVIEW, - WSS_PREPARE_BUILDER_SUMMARY, - WSS_SEND_BUILDER_SUMMARY, - - WSS_PREPARE_TASKINFO, - WSS_SEND_ARTIFACT_INFO, - - WSS_PREPARE_EVENTINFO, - WSS_SEND_EVENTINFO, -} ws_state; - typedef struct sai_builder { sais_t c; } saib_t; @@ -75,32 +59,6 @@ enum { SAIM_SPECIFIC_TASK, }; -typedef struct saiw_scheduled { - lws_dll2_t list; - - char task_uuid[65]; - - sai_task_t *one_task; /* only for browser */ - const sai_event_t *one_event; - - lws_dll2_t *walk; - - lws_dll2_owner_t owner; - - struct lwsac *ac; - struct lwsac *query_ac; /* taskinfo event only */ - - ws_state action; - int task_index; - - uint8_t ovstate; /* SOS_ substate when doing overview */ - - uint8_t subsequent:1; /* for individual JSON */ - uint8_t ov_db_done:1; /* for individual JSON */ - uint8_t logsub:1; /* for individual JSON */ - -} saiw_scheduled_t; - struct pss { struct vhd *vhd; struct lws *wsi; @@ -123,6 +81,7 @@ struct pss { sqlite3_blob *blob_artifact; lws_dll2_owner_t logs_owner; + lws_sorted_usec_list_t sul_logcache; lws_struct_args_t a; union { @@ -155,8 +114,6 @@ struct pss { uint64_t artifact_offset; uint64_t artifact_length; - ws_state send_state; - unsigned int spa_failed:1; unsigned int dry:1; unsigned int frag:1; @@ -194,7 +151,6 @@ struct vhd { lws_dll2_owner_t sqlite3_cache; /* sais_sqlite_cache_t */ lws_dll2_owner_t tasklog_cache; - lws_sorted_usec_list_t sul_logcache; }; typedef struct saiw_websrv { @@ -264,18 +220,9 @@ saiw_get_blob(struct vhd *vhd, const char *url, sqlite3 **pdb, int saiw_browsers_task_state_change(struct vhd *vhd, const char *task_uuid); -saiw_scheduled_t * -saiw_alloc_sched(struct pss *pss, ws_state action); void -saiw_dealloc_sched(saiw_scheduled_t *sch); - -int -saiw_sched_destroy(struct lws_dll2 *d, void *user); - - -void -saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, +saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(struct vhd *vhd, const void *buf, size_t len, enum lws_write_protocol flags); void @@ -284,3 +231,13 @@ saiw_browser_state_changed(struct pss *pss, int established); void saiw_update_viewer_count(struct vhd *vhd); +int +saiw_broadcast_logs_batch(struct vhd *vhd, struct pss *pss); + +int +saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss); + +int +saiw_browser_broadcast_queue_builders(struct vhd *vhd, struct pss *pss); + + diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 0d0d625..b6d0804 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -135,6 +135,43 @@ enum sai_overview_state { SOS_TASKS, }; +int +saiw_ws_browser_queue_REQUIRES_LWS_PRE(struct pss *pss, const void *buf, + size_t len, enum lws_write_protocol flags) +{ + int *pi = (int *)((const char *)buf - sizeof(int)), r = 0; + + *pi = (int)flags; + + if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) < 0) { + lwsl_wsi_err(pss->wsi, "unable to buflist_append"); /* still ask to drain */ + r = 1; + } + + lws_callback_on_writable(pss->wsi); + + return r; +} + +/* + * This allows other parts of sai-web to queue a raw buffer to be sent to + * all connected browsers, eg, for load reports. + * + * The flags are lws_write() flags. + */ +void +saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(struct vhd *vhd, const void *buf, + size_t len, enum lws_write_protocol flags) +{ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, buf, len, flags); + + } lws_end_foreach_dll(p); +} + + int sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name) @@ -179,44 +216,11 @@ saiw_subs_request_writeable(struct vhd *vhd, const char *task_uuid) return 0; } -saiw_scheduled_t * -saiw_alloc_sched(struct pss *pss, ws_state action) -{ - saiw_scheduled_t *sch = malloc(sizeof(*sch)); - - if (sch) { - memset(sch, 0, sizeof(*sch)); - sch->action = action; - lws_dll2_add_tail(&sch->list, &pss->sched); - lws_callback_on_writable(pss->wsi); - } - - return sch; -} - -void -saiw_dealloc_sched(saiw_scheduled_t *sch) -{ - if (!sch) - return; - - lws_dll2_remove(&sch->list); - - lwsac_free(&sch->ac); - lwsac_free(&sch->query_ac); - - free(sch); -} - static int saiw_pss_schedule_eventinfo(struct pss *pss, const char *event_uuid) { - saiw_scheduled_t *sch = saiw_alloc_sched(pss, WSS_PREPARE_OVERVIEW); - char qu[180], esc[66], esc2[96]; - int n; - - if (!sch) - return -1; +// char qu[180], esc[66], esc2[96]; +// int n; /* * This pss may be locked to a specific event @@ -231,7 +235,7 @@ saiw_pss_schedule_eventinfo(struct pss *pss, const char *event_uuid) * * Just collect the event struct into pss->query_owner to dump */ - +#if 0 lws_sql_purify(esc, event_uuid, sizeof(esc)); if (pss->specific_project[0]) { @@ -244,15 +248,16 @@ saiw_pss_schedule_eventinfo(struct pss *pss, const char *event_uuid) &sch->owner, &sch->ac, 0, 1); if (n < 0 || !sch->owner.head) goto bail; - - sch->ov_db_done = 1; - // lwsl_warn("%s: doing WSS_PREPARE_BUILDER_SUMMARY\n", __func__); - saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); +#endif + saiw_browser_queue_overview(pss->vhd, pss); + saiw_browser_broadcast_queue_builders(pss->vhd, pss); return 0; bail: - saiw_dealloc_sched(sch); + saiw_browser_queue_overview(pss->vhd, pss); + saiw_browser_broadcast_queue_builders(pss->vhd, pss); + return 1; } @@ -261,15 +266,21 @@ bail: static int saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) { - saiw_scheduled_t *sch = saiw_alloc_sched(pss, WSS_PREPARE_TASKINFO); - char qu[192], esc[66], event_uuid[33], esc2[96]; + char qu[192], event_uuid[33], esc2[96], buf[4096 + LWS_PRE], + *start = buf + LWS_PRE, *p = start, *end = buf + sizeof(buf); + const sai_event_t *one_event = NULL; + sai_browse_taskreply_t task_reply; + struct lwsac *query_ac = NULL; + sai_task_t *one_task = NULL; + lws_struct_serialize_t *js; + char esc[256], filt[128]; + lws_dll2_owner_t owner; sqlite3 *pdb = NULL; lws_dll2_owner_t o; sai_task_t *pt; - int n, m; - - if (!sch) - return -1; + char fi = 1; + int m, n; + size_t w; sai_task_uuid_to_event_uuid(event_uuid, task_uuid); @@ -290,11 +301,8 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) /* open the event-specific database object */ if (sai_event_db_ensure_open(pss->vhd->context, &pss->vhd->sqlite3_cache, - pss->vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { - /* no longer exists, nothing to do */ - saiw_dealloc_sched(sch); + pss->vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) return 0; - } /* * get the related task object into its own ac... there might @@ -305,17 +313,17 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) lws_sql_purify(esc, task_uuid, sizeof(esc)); lws_snprintf(qu, sizeof(qu), " and uuid='%s'", esc); n = lws_struct_sq3_deserialize(pdb, qu, NULL, lsm_schema_sq3_map_task, - &o, &sch->query_ac, 0, 1); + &o, &query_ac, 0, 1); sai_event_db_close(&pss->vhd->sqlite3_cache, &pdb); if (n < 0 || !o.head) goto bail; pt = lws_container_of(o.head, sai_task_t, list); - sch->one_task = pt; + one_task = pt; /* let the pss take over the task info ac and schedule sending */ - lws_dll2_remove((struct lws_dll2 *)&sch->one_task->list); + lws_dll2_remove((struct lws_dll2 *)&one_task->list); /* * let's also get the event object the task relates to into @@ -348,25 +356,157 @@ 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); + &query_ac, 0, 1); 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; + one_event = NULL; else - sch->one_event = lws_container_of(o.head, sai_event_t, list); + one_event = lws_container_of(o.head, sai_event_t, list); + + memset(&task_reply, 0, sizeof(task_reply)); + + /* + * We're sending a browser the specific task info that he + * asked for. + * + * We already got the task struct out of the db in .one_task + * (all in .query_ac)... we're responsible for destroying it + * when we go out of scope... + */ + + task_reply.event = one_event; + task_reply.task = one_task; + one_task->rebuildable = (one_task->state == SAIES_FAIL || + one_task->state == SAIES_CANCELLED) && + (lws_now_secs() - (one_task->started + + (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; + lws_strncpy(task_reply.auth_user, pss->auth_user, sizeof(task_reply.auth_user)); + + js = lws_struct_json_serialize_create(lsm_schema_json_map_taskreply, + LWS_ARRAY_SIZE(lsm_schema_json_map_taskreply), + 0, &task_reply); + if (!js) { + lwsl_warn("%s: couldn't create\n", __func__); + goto bail; + } + + do { + n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w); + + if (lws_ptr_diff_size_t(end, (uint8_t *)p) < 512) { + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(pss->vhd, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, fi, 0)); + p = start; + fi = 0; + } + + } while (n == LSJS_RESULT_CONTINUE); + + lws_struct_json_serialize_destroy(&js); + + /* + * Let's also try to fetch any artifacts into pss->aft_owner... + * no db or no artifacts can also be a normal situation... + */ + + if (one_task) { + + sai_task_uuid_to_event_uuid(event_uuid, one_task->uuid); + + lws_dll2_owner_clear(&owner); + if (!sai_event_db_ensure_open(pss->vhd->context, &pss->vhd->sqlite3_cache, + pss->vhd->sqlite3_path_lhs, event_uuid, + 0, &pdb)) { + + lws_snprintf(filt, sizeof(filt), " and (task_uuid == '%s')", + one_task->uuid); + + if (lws_struct_sq3_deserialize(pdb, filt, NULL, + lsm_schema_sq3_map_artifact, + &owner, + &query_ac, 0, 10)) + lwsl_err("%s: get afcts failed\n", __func__); + + sai_event_db_close(&pss->vhd->sqlite3_cache, &pdb); + } + } + + if (n == LSJS_RESULT_ERROR) { + lwsl_notice("%s: taskinfo: error generating json\n", __func__); + goto bail; + } + p += w; + + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, fi, 1)); + + /* does he want to subscribe to logs? */ + if (logsub && one_task && !pss->subs_list.owner) { + strcpy(pss->sub_task_uuid, one_task->uuid); + lws_dll2_add_head(&pss->subs_list, &pss->vhd->subs_owner); + pss->sub_timestamp = pss->initial_log_timestamp; /* where we got up to */ + saiw_broadcast_logs_batch(pss->vhd, pss); + } + + saiw_browser_broadcast_queue_builders(pss->vhd, pss); + + if (owner.head) { + sai_artifact_t *aft = (sai_artifact_t *)owner.head; + + p = start; + fi = 1; + + lwsl_info("%s: WSS_SEND_ARTIFACT_INFO: consuming artifact\n", __func__); + + lws_dll2_remove(&aft->list); + + /* we don't want to disclose this to browsers */ + aft->artifact_up_nonce[0] = '\0'; - sch->logsub = !!logsub; + js = lws_struct_json_serialize_create(lsm_schema_json_map_artifact, + LWS_ARRAY_SIZE(lsm_schema_json_map_artifact), + 0, aft); + if (!js) { + lwsl_err("%s ----------------- failed to render artifact json\n", __func__); + goto bail; + } + + do { + n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w); + if (n == LSJS_RESULT_ERROR) { + lws_struct_json_serialize_destroy(&js); + lwsl_notice("%s: taskinfo: ---------- error generating json\n", __func__); + goto bail; + } + p += w; + if (lws_ptr_diff_size_t(end, p) < 512) { + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(pss->vhd, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, fi, 0)); + p = start; + fi = 0; + } - saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); + } while (n == LSJS_RESULT_CONTINUE); + + lws_struct_json_serialize_destroy(&js); + } + + lwsac_free(&query_ac); return 0; bail: - saiw_dealloc_sched(sch); + lwsac_free(&query_ac); + return 1; } @@ -474,8 +614,8 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, if (ti->js_api_version) pss->js_api_version = ti->js_api_version; - saiw_alloc_sched(pss, WSS_PREPARE_OVERVIEW); - saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); + saiw_browser_broadcast_queue_builders(pss->vhd, pss); + saiw_browser_queue_overview(pss->vhd, pss); break; } @@ -638,386 +778,268 @@ soft_error: return 0; } +static void +saiw_retry_logs(lws_sorted_usec_list_t *sul) +{ + struct pss *pss = lws_container_of(sul, struct pss, sul_logcache); -/* - * We're sending something on a browser ws connection. Returning nonzero from - * here drops the connection, necessary if we fail partway through a message - * but undesirable if a browser tab will keep reconnecting and asking for the - * same, no-longer-existant thing. - */ + saiw_broadcast_logs_batch(pss->vhd, pss); +} int -saiw_ws_json_tx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) +saiw_broadcast_logs_batch(struct vhd *vhd, struct pss *pss) { - uint8_t *start = buf + LWS_PRE, *p = start, *end = p + bl - LWS_PRE - 1; - int n, flags = LWS_WRITE_TEXT, first = 0, iu, endo; - char esc[256], esc1[33], filt[128]; - sai_browse_taskreply_t task_reply; - struct lwsac *task_ac = NULL; - lws_dll2_owner_t task_owner; - lws_struct_serialize_t *js; - saiw_scheduled_t *sch; char event_uuid[33]; - sqlite3 *pdb = NULL; - sai_event_t *e; - sai_task_t *t; - char any, lg; - size_t w; - -again: - - start = buf + LWS_PRE; - p = start; - end = p + bl - LWS_PRE - 1; - flags = LWS_WRITE_TEXT; - first = 0; - lg = 0; - endo = 0; - - // lwsl_notice("%s: send_state %d, pss %p, wsi %p\n", __func__, - // pss->send_state, pss, pss->wsi); - - sch = NULL; - if (pss->sched.head) - sch = lws_container_of(pss->sched.head, saiw_scheduled_t, list); - - switch (pss->send_state) { - case WSS_IDLE1: - - /* - * Anything from a task log he's subscribed to? - * - * If so, let's prioritize that first... - */ - - if ((!pss->sched.head || !pss->toggle_favour_sch) && - pss->subs_list.owner) { - - sch = NULL; - - /* - * For efficiency, let's try to grab the next 100 at - * once from sqlite and work our way through sending - * them - */ - - if (pss->log_cache_index == pss->log_cache_size) { - int sr; - - sai_task_uuid_to_event_uuid(event_uuid, - pss->sub_task_uuid); - lws_dll2_owner_clear(&task_owner); - lwsac_free(&pss->logs_ac); - - lws_snprintf(esc, sizeof(esc), - "and task_uuid='%s' and timestamp > %llu", - pss->sub_task_uuid, - (unsigned long long)pss->sub_timestamp); - - lwsl_info("%s: collecting logs %s\n", - __func__, esc); - - if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, - vhd->sqlite3_path_lhs, event_uuid, 0, - &pdb)) { - lwsl_notice("%s: unable to open event-specific database\n", - __func__); - - return 0; - } + if (!pss->subs_list.owner) + return 0; - sr = lws_struct_sq3_deserialize(pdb, esc, - "uid,timestamp ", - lsm_schema_sq3_map_log, - &pss->logs_owner, - &pss->logs_ac, 0, 100); + /* + * For efficiency, let's try to grab the next 100 at + * once from sqlite and work our way through sending + * them + */ - sai_event_db_close(&vhd->sqlite3_cache, &pdb); + //if (pss->log_cache_index == pss->log_cache_size) + { + sqlite3 *pdb = NULL; + char esc[256]; + int sr; - if (sr) { + sai_task_uuid_to_event_uuid(event_uuid, pss->sub_task_uuid); - lwsl_err("%s: subs failed\n", __func__); + lwsac_free(&pss->logs_ac); - return 0; - } + lws_snprintf(esc, sizeof(esc), + "and task_uuid='%s' and timestamp > %llu", + pss->sub_task_uuid, + (unsigned long long)pss->sub_timestamp); - pss->log_cache_index = 0; - pss->log_cache_size = (int)pss->logs_owner.count; - } + // lwsl_notice("%s: collecting logs %s\n", __func__, esc); - if (pss->log_cache_index < pss->log_cache_size) { - sai_log_t *log = lws_container_of( - pss->logs_owner.head, - sai_log_t, list); + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, + 0, &pdb)) { + lwsl_notice("%s: unable to open event-specific database\n", + __func__); - lws_dll2_remove(&log->list); - pss->log_cache_index++; + return 0; + } - /* - * Turn it back into JSON so we can give it to - * the browser - */ + sr = lws_struct_sq3_deserialize(pdb, esc, + "uid,timestamp ", + lsm_schema_sq3_map_log, + &pss->logs_owner, + &pss->logs_ac, 0, 50); - js = lws_struct_json_serialize_create( - lsm_schema_json_map_log, 1, 0, log); - if (!js) { - lwsl_notice("%s: json ser fail\n", __func__); - return 0; - } + sai_event_db_close(&vhd->sqlite3_cache, &pdb); - n = (int)lws_struct_json_serialize(js, p, - lws_ptr_diff_size_t(end, p), &w); - lws_struct_json_serialize_destroy(&js); - if (n == LSJS_RESULT_ERROR) { - lwsl_notice("%s: json ser error\n", __func__); - return 0; - } + if (sr) { - p += w; - first = 1; - lg = 1; - pss->toggle_favour_sch = 1; - - /* - * Record that this was the most recent log we - * saw so far - */ - pss->sub_timestamp = log->timestamp; - goto send_it; - } - } - - /* - * Stay in this state if we're in the middle of a - * multi-fragment message - */ - if (lws_ws_sending_multifragment(pss->wsi)) { - lws_callback_on_writable(pss->wsi); + lwsl_err("%s: subs failed\n", __func__); return 0; } - /* fallthru */ + pss->log_cache_index = 0; + pss->log_cache_size = (int)pss->logs_owner.count; + } - case WSS_IDLE2: + while (pss->log_cache_index++ < pss->log_cache_size) { + sai_log_t *log = lws_container_of(pss->logs_owner.head, + sai_log_t, list); + lws_struct_serialize_t *js; + char buf[1200 + LWS_PRE]; + char fi = 1; + int n; - pss->send_state = WSS_IDLE2; + lws_dll2_remove(&log->list); /* - * Send anything waiting on broadcast_raw buflist first + * Turn it back into JSON so we can give it to + * the browser */ - if (pss->raw_tx) { - /* - * 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. - */ - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&pss->raw_tx), depi = *pi; - char som, eom, rb[1200]; - int used, final = 1; - size_t fsl = lws_buflist_next_segment_len(&pss->raw_tx, NULL); - - /* this is the only buflist user on pss->raw_tx */ - used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, sizeof(rb), &som, &eom); - if (!used) - return 0; - if (used < (int)fsl || (depi & LWS_WRITE_NO_FIN)) - final = 0; - - if (lws_write(pss->wsi, (uint8_t *)rb + ((size_t)som * sizeof(int)), - (size_t)used - ((size_t)som * sizeof(int)), - (lws_ws_sending_multifragment(pss->wsi) ? LWS_WRITE_CONTINUATION : LWS_WRITE_TEXT) | - (!final * LWS_WRITE_NO_FIN)) < 0) { - lwsl_wsi_err(pss->wsi, "attempt to write %d failed", (int)used - (int)sizeof(int)); - - return -1; - } - - if (lws_buflist_next_segment_len(&pss->raw_tx, NULL)) - lws_callback_on_writable(pss->wsi); - - if (!lws_ws_sending_multifragment(pss->wsi)) - pss->send_state = WSS_IDLE1; - + js = lws_struct_json_serialize_create(lsm_schema_json_map_log, + 1, 0, log); + if (!js) { + lwsl_notice("%s: json ser fail\n", __func__); return 0; } - /* - * Stay in this state if we're in the middle of a - * multi-fragment message, otherwise do whatever the - * sch proposes - */ - - if (lws_ws_sending_multifragment(pss->wsi) || - !sch) - return 0; - - /* switch to the pending sch */ - - pss->toggle_favour_sch = 0; - pss->send_state = sch->action; - goto again; - - case WSS_PREPARE_OVERVIEW: + do { + size_t w; + n = lws_struct_json_serialize(js, (uint8_t *)buf + LWS_PRE, + sizeof(buf) - LWS_PRE, &w); - if (!sch) /* coverity */ - goto no_sch; + if (n != LSJS_RESULT_CONTINUE) + lws_struct_json_serialize_destroy(&js); + if (n == LSJS_RESULT_ERROR) + return 1; - filt[0] = '\0'; - if (pss->specific_project[0]) { - lws_sql_purify(esc, pss->specific_project, sizeof(esc) - 1); - lws_snprintf(filt, sizeof(filt), " and repo_name=\"%s\"", esc); - } - if (!pss->authorized) - lws_snprintf(filt + strlen(filt), sizeof(filt) - strlen(filt), " and sec=0"); + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, buf + LWS_PRE, w, + lws_write_ws_flags(LWS_WRITE_TEXT, + fi, n == LSJS_RESULT_FINISH)); - pss->wants_event_updates = 1; - if (!sch->ov_db_done && lws_struct_sq3_deserialize(vhd->pdb, - filt[0] ? filt : NULL, "created ", - lsm_schema_sq3_map_event, &sch->owner, - &sch->ac, 0, -8)) { - lwsl_notice("%s: OVERVIEW 2 failed\n", __func__); + fi = 0; + pss->sub_timestamp = log->timestamp; + } while (n != LSJS_RESULT_FINISH); + } - pss->send_state = WSS_IDLE1; - saiw_dealloc_sched(sch); + lwsac_free(&pss->logs_ac); - return 0; - } + lws_sul_schedule(vhd->context, 0, &pss->sul_logcache, + saiw_retry_logs, + pss->log_cache_size == 50 ? 500 : 250 * LWS_US_PER_MS); - /* - * we get zero or more sai_event_t laid out in pss->query_ac, - * and listed in pss->query_owner - */ + return 0; +} - lwsl_debug("%s: WSS_PREPARE_OVERVIEW: %d results %p\n", - __func__, sch->owner.count, sch->ac); - - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), - "{\"schema\":\"sai.warmcat.com.overview\"," - " \"api_version\":%u," - " \"alang\":\"%s\"," - " \"authorized\": %d," - " \"auth_secs\": %ld," - " \"auth_user\": \"%s\"," - "\"overview\":[", - SAIW_API_VERSION, - lws_json_purify(esc, pss->alang, sizeof(esc) - 1, &iu), - pss->authorized, pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0, - lws_json_purify(esc1, pss->auth_user, sizeof(esc1) - 1, &iu) - ); +int +saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss) +{ + char buf[4096 + LWS_PRE], *start = buf + LWS_PRE, *p = start, + *end = buf + sizeof(buf); + char esc[256], esc1[33], filt[128], subsequent; + struct lwsac *task_ac = NULL, *ac = NULL; + lws_dll2_owner_t task_owner, owner; + unsigned int task_index = 0; + lws_struct_serialize_t *js; + sqlite3 *pdb = NULL; + lws_dll2_t *walk; + sai_task_t *t; + int n, iu; + size_t w; - /* - * "authorized" here is used to decide whether to render the - * additional controls clientside. The events the controls - * cause if used are separately checked for coming from an - * authorized pss when they are received. - * - * If you're not authorized, you're only going to see events - * that have sec=0. Otherwise you can see all events. - */ + filt[0] = '\0'; + esc[0] = '\0'; + n = -8; - if (pss->specificity) - sch->walk = lws_dll2_get_head(&sch->owner); - else - sch->walk = lws_dll2_get_tail(&sch->owner); - sch->subsequent = 0; - first = 1; + if (pss->specific_project[0]) { + lws_sql_purify(esc, pss->specific_project, sizeof(esc) - 1); + lws_snprintf(filt, sizeof(filt), " and repo_name=\"%s\"", esc); + n = -1; + } + if (!pss->authorized) + lws_snprintf(filt + strlen(filt), sizeof(filt) - strlen(filt), " and sec=0"); - pss->send_state = WSS_SEND_OVERVIEW; + pss->wants_event_updates = 1; + if (lws_struct_sq3_deserialize(vhd->pdb, filt[0] ? filt : NULL, + "created ", lsm_schema_sq3_map_event, + &owner, &ac, 0, n)) { + lwsl_notice("%s: OVERVIEW 2 failed\n", __func__); - if (!sch->owner.count) - goto so_finish; + return 0; + } - /* fallthru */ + /* + * we get zero or more sai_event_t laid out in pss->query_ac, + * and listed in pss->query_owner + */ - case WSS_SEND_OVERVIEW: + p += (size_t)lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), + "{\"schema\":\"sai.warmcat.com.overview\"," + " \"api_version\":%u," + " \"alang\":\"%s\"," + " \"authorized\": %d," + " \"auth_secs\": %ld," + " \"auth_user\": \"%s\"," + "\"overview\":[", SAIW_API_VERSION, + lws_json_purify(esc, pss->alang, sizeof(esc) - 1, &iu), + pss->authorized, pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0, + lws_json_purify(esc1, pss->auth_user, sizeof(esc1) - 1, &iu) + ); + + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, 1, 0)); + p = start; - if (!sch) /* coverity */ - goto no_sch; - if (sch->ovstate == SOS_TASKS) - goto enum_tasks; + /* + * "authorized" here is used to decide whether to render the + * additional controls clientside. The events the controls + * cause if used are separately checked for coming from an + * authorized pss when they are received. + * + * If you're not authorized, you're only going to see events + * that have sec=0. Otherwise you can see all events. + */ - any = 0; - while (end - p > 2048 && sch->walk && - pss->send_state == WSS_SEND_OVERVIEW) { + if (pss->specificity) + walk = lws_dll2_get_head(&owner); + else + walk = lws_dll2_get_tail(&owner); - e = lws_container_of(sch->walk, sai_event_t, list); + subsequent = 0; - if (pss->specificity) { - lwsl_debug("%s: Specificity: e->hash: %s, " - "e->ref: '%s', pss->specific_ref: '%s'\n", - __func__, e->hash, e->ref, - pss->specific_ref); + if (!owner.count) /* nothing to do */ + goto so_finish; - if (!strcmp(pss->specific_ref, "refs/heads/master") && - !strcmp(e->ref, "refs/heads/main")) { - // lwsl_notice("master->main\n"); - any = 1; - } else { + while (walk) { + sai_event_t *e = lws_container_of(walk, sai_event_t, list); - if (strcmp(e->hash, pss->specific_ref) && - strcmp(e->ref, pss->specific_ref)) { - sch->walk = sch->walk->next; - continue; - } - //lwsl_notice("%s: match\n", __func__); - any = 1; + if (pss->specificity) { + if (!strcmp(pss->specific_ref, "refs/heads/master") && + !strcmp(e->ref, "refs/heads/main")) + ; // any = 1; + else { + if (strcmp(e->hash, pss->specific_ref) && + strcmp(e->ref, pss->specific_ref)) { + walk = walk->next; + continue; } + // any = 1; } + } - js = lws_struct_json_serialize_create( - lsm_schema_json_map_event, - LWS_ARRAY_SIZE(lsm_schema_json_map_event), 0, e); - if (!js) { - lwsl_err("%s: json ser fail\n", __func__); - return 1; - } - if (sch->subsequent) - *p++ = ','; - sch->subsequent = 1; + js = lws_struct_json_serialize_create( + lsm_schema_json_map_event, + LWS_ARRAY_SIZE(lsm_schema_json_map_event), 0, e); + if (!js) { + lwsl_err("%s: json ser fail\n", __func__); + return 1; + } + if (subsequent) + *p++ = ','; + subsequent = 1; - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "{\"e\":"); + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "{\"e\":"); - n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w); - lws_struct_json_serialize_destroy(&js); - switch (n) { - case LSJS_RESULT_ERROR: - pss->send_state = WSS_IDLE1; - saiw_dealloc_sched(sch); - lwsl_err("%s: json ser error\n", __func__); - return 1; + if (lws_ptr_diff_size_t(end, p) < 256) { + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0)); + p = start; + } - case LSJS_RESULT_FINISH: - case LSJS_RESULT_CONTINUE: - p += w; - sch->ovstate = SOS_TASKS; - sch->task_index = 0; - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), ", \"t\":["); - goto enum_tasks; - } + n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w); + lws_struct_json_serialize_destroy(&js); + switch (n) { + case LSJS_RESULT_ERROR: + lwsl_err("%s: json ser error\n", __func__); + return 1; + + case LSJS_RESULT_FINISH: + case LSJS_RESULT_CONTINUE: + p += w; + task_index = 0; + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), ", \"t\":["); + break; } - if (!any) { - pss->send_state = WSS_IDLE1; - saiw_dealloc_sched(sch); - return 0; + if (lws_ptr_diff_size_t(end, p) < 2560) { + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0)); + p = start; } - break; -enum_tasks: /* - * Enumerate the tasks associated with this event... we will - * come back here as often as needed to dump all the tasks + * Enumerate the tasks associated with this event... */ - e = lws_container_of(sch->walk, sai_event_t, list); + e = lws_container_of(walk, sai_event_t, list); lws_dll2_owner_clear(&task_owner); do { @@ -1034,7 +1056,7 @@ enum_tasks: lws_dll2_owner_clear(&task_owner); if (lws_struct_sq3_deserialize(pdb, NULL, NULL, lsm_schema_sq3_map_task, &task_owner, - &task_ac, sch->task_index, 1)) { + &task_ac, (int)task_index, 1)) { lwsl_err("%s: OVERVIEW 1 failed\n", __func__); sai_event_db_close(&vhd->sqlite3_cache, &pdb); @@ -1045,7 +1067,7 @@ enum_tasks: if (!task_owner.count) break; - if (sch->task_index) + if (task_index) *p++ = ','; /* @@ -1074,316 +1096,116 @@ enum_tasks: lsm_schema_json_map_task, LWS_ARRAY_SIZE(lsm_schema_json_map_task), 0, t); - n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w); + t->build[0] = '\0'; + n = (int)lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w); lws_struct_json_serialize_destroy(&js); lwsac_free(&task_ac); p += w; - sch->task_index++; - } while ((end - p > 2048) && task_owner.count); + if (lws_ptr_diff_size_t(end, p) < 2560) { + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, 0, 0)); + p = start; + } - if (task_owner.count) - /* may be more left to do */ - break; + task_index++; + } while (1); /* none left to do, go back up a level */ p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - sch->ovstate = SOS_EVENT; if (pss->specificity) - sch->walk = sch->walk->next; + walk = walk->next; else - sch->walk = sch->walk->prev; - if (!sch->walk || pss->specificity) { - while (sch->walk) - sch->walk = sch->walk->next; - goto so_finish; - } - break; - -so_finish: - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - pss->send_state = WSS_IDLE1; - endo = 1; - break; - - case WSS_PREPARE_BUILDER_SUMMARY: - - if (!sch) /* coverity */ - goto no_sch; - - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), - "{\"schema\":\"com.warmcat.sai.builders\"," - " \"alang\":\"%s\"," - " \"authorized\":%d," - " \"auth_secs\":%ld," - " \"auth_user\": \"%s\"," - " \"builders\":[", - lws_sql_purify(esc, pss->alang, sizeof(esc) - 1), - pss->authorized, pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0, - lws_json_purify(esc1, pss->auth_user, sizeof(esc1) - 1, &iu)); - - if (vhd && vhd->builders) { - // lwsac_reference(vhd->builders); - sch->walk = lws_dll2_get_head(&vhd->builders_owner); - - /* HEAD of the owner list must be inside the vhd->builders ac */ - // if (sch->walk && lwsac_assert_valid(vhd->builders, sch->walk, sizeof(sai_plat_t))) - // break; - } else { - lwsl_notice("%s: BUILDER_SUMMARY: can't start walk\n", __func__); - sch->walk = 0; - } - -// sch->walk = 0; - - sch->subsequent = 0; - pss->send_state = WSS_SEND_BUILDER_SUMMARY; - first = 1; - - /* fallthru */ - - case WSS_SEND_BUILDER_SUMMARY: - - if (!sch) /* coverity */ - goto no_sch; - - if (!sch->walk) - goto b_finish; - - /* - * We're going to send the browser some JSON about all the - * builders / platforms we feel are connected to us - */ - - while (end - p > 512 && sch->walk && - pss->send_state == WSS_SEND_BUILDER_SUMMARY) { - - /* every builder must be inside the vhd->builders ac */ - //if (lwsac_assert_valid(vhd->builders, sch->walk, sizeof(sai_plat_t))) - // break; - - sai_plat_t *b = lws_container_of(sch->walk, sai_plat_t, - sai_plat_list); - - js = lws_struct_json_serialize_create( - lsm_schema_map_plat_simple, - LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), - 0, b); - if (!js) { - lwsac_unreference(&vhd->builders); - return 1; - } - if (sch->subsequent) - *p++ = ','; - sch->subsequent = 1; - - switch (lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w)) { - case LSJS_RESULT_ERROR: - lws_struct_json_serialize_destroy(&js); - pss->send_state = WSS_IDLE1; - saiw_dealloc_sched(sch); - return 1; - - case LSJS_RESULT_FINISH: - case LSJS_RESULT_CONTINUE: - p += w; - lws_struct_json_serialize_destroy(&js); - sch->walk = sch->walk->next; - if (!sch->walk) - goto b_finish; - break; - } - } - break; -b_finish: - p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - - // lwsac_unreference(&vhd->builders); - endo = 1; - break; + walk = walk->prev; + if (walk && !pss->specificity) + continue; + } - case WSS_PREPARE_TASKINFO: +so_finish: + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - if (!sch) /* coverity */ - goto no_sch; + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, 0, 1)); - /* - * We're sending a browser the specific task info that he - * asked for. - * - * We already got the task struct out of the db in .one_task - * (all in .query_ac)... we're responsible for destroying it - * when we go out of scope... - */ + return 0; +} - 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; - lws_strncpy(task_reply.auth_user, pss->auth_user, - sizeof(task_reply.auth_user)); - - js = lws_struct_json_serialize_create(lsm_schema_json_map_taskreply, - LWS_ARRAY_SIZE(lsm_schema_json_map_taskreply), - 0, &task_reply); +int +saiw_browser_broadcast_queue_builders(struct vhd *vhd, struct pss *pss) +{ + char buf[4096 + LWS_PRE], *start = buf + LWS_PRE, *p = start, + *end = buf + sizeof(buf); + lws_struct_serialize_t *js; + char esc[256], esc1[33]; + lws_dll2_t *walk = NULL; + char fi = 1, subsequent; + size_t w; + int iu; + + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), + "{\"schema\":\"com.warmcat.sai.builders\"," + " \"alang\":\"%s\"," + " \"authorized\":%d," + " \"auth_secs\":%ld," + " \"auth_user\": \"%s\"," + " \"builders\":[", + lws_sql_purify(esc, pss->alang, sizeof(esc) - 1), + pss->authorized, pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0, + lws_json_purify(esc1, pss->auth_user, sizeof(esc1) - 1, &iu)); + + if (vhd && vhd->builders) + walk = lws_dll2_get_head(&vhd->builders_owner); + + subsequent = 0; + + while (walk) { + sai_plat_t *b = lws_container_of(walk, sai_plat_t, sai_plat_list); + + js = lws_struct_json_serialize_create( + lsm_schema_map_plat_simple, + LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), + 0, b); if (!js) { - saiw_dealloc_sched(sch); - lwsl_warn("%s: couldn't create\n", __func__); + lwsac_unreference(&vhd->builders); return 1; } + if (subsequent) + *p++ = ','; + subsequent = 1; - n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w); - lws_struct_json_serialize_destroy(&js); - - /* - * Let's also try to fetch any artifacts into pss->aft_owner... - * no db or no artifacts can also be a normal situation... - */ - - if (sch->one_task) { - - sai_task_uuid_to_event_uuid(event_uuid, - sch->one_task->uuid); - - //lwsl_debug("%s: ---------------- event uuid '%s'\n", - // __func__, event_uuid); - - lws_dll2_owner_clear(&sch->owner); - if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, - vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { - - lws_snprintf(filt, sizeof(filt), " and (task_uuid == '%s')", - sch->one_task->uuid); - - // if (!pss->authorized) - // lws_snprintf(filt + strlen(filt), sizeof(filt) - strlen(filt), " and sec=0"); - - // lwsl_debug("%s: ---------------- %s\n", __func__, filt); - - if (lws_struct_sq3_deserialize(pdb, filt, NULL, - lsm_schema_sq3_map_artifact, - &sch->owner, - &sch->ac, 0, 10)) { - lwsl_err("%s: get afcts failed\n", __func__); - } - sai_event_db_close(&vhd->sqlite3_cache, &pdb); - } - } - - first = 1; - sch->walk = NULL; - pss->send_state = WSS_SEND_ARTIFACT_INFO; - if (!sch->owner.head) { - // lwsl_debug("%s: ---------------- no artifacts\n", __func__); - /* there's no artifact stuff to do */ - endo = 1; - } else - lwsl_debug("%s: WSS_PREPARE_TASKINFO: planning on artifacts\n", __func__); - // sch->one_task = NULL; - if (n == LSJS_RESULT_ERROR) { - saiw_dealloc_sched(sch); - lwsl_notice("%s: taskinfo: error generating json\n", __func__); + switch (lws_struct_json_serialize(js, (uint8_t *)p, lws_ptr_diff_size_t(end, p), &w)) { + case LSJS_RESULT_ERROR: + lws_struct_json_serialize_destroy(&js); return 1; - } - p += w; - if (!lws_ptr_diff(p, start)) { - saiw_dealloc_sched(sch); - pss->send_state = WSS_IDLE1; - lwsl_notice("%s: taskinfo: empty json\n", __func__); - return 0; - } - -// sai_dump_stderr((const char *)start, lws_ptr_diff_size_t(p, start)); - break; - - case WSS_SEND_ARTIFACT_INFO: - - if (!sch) /* coverity */ - goto no_sch; - if (sch->owner.head) { - sai_artifact_t *aft = (sai_artifact_t *)sch->owner.head; - - lwsl_info("%s: WSS_SEND_ARTIFACT_INFO: consuming artifact\n", __func__); - - lws_dll2_remove(&aft->list); - - /* we don't want to disclose this to browsers */ - aft->artifact_up_nonce[0] = '\0'; - - js = lws_struct_json_serialize_create(lsm_schema_json_map_artifact, - LWS_ARRAY_SIZE(lsm_schema_json_map_artifact), - 0, aft); - if (!js) { - saiw_dealloc_sched(sch); - lwsl_err("%s ----------------- failed to render artifact json\n", __func__); - return 1; - } - - n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w); + case LSJS_RESULT_FINISH: lws_struct_json_serialize_destroy(&js); - if (n == LSJS_RESULT_ERROR) { - saiw_dealloc_sched(sch); - lwsl_notice("%s: taskinfo: ---------- error generating json\n", __func__); - return 1; - } - first = 1; + /* fallthru */ + case LSJS_RESULT_CONTINUE: p += w; - // lwsl_warn("%s: --------------------- %.*s\n", __func__, (int)w, start); + walk = walk->next; + break; } - if (!sch->owner.head) - endo = 1; - break; - - default: - lwsl_err("%s: pss state %d\n", __func__, pss->send_state); - return 0; - } - -send_it: - - flags = lws_write_ws_flags(LWS_WRITE_TEXT, first, endo || lg || (sch && !sch->walk)); - - if (lg || - endo || - (pss->send_state == WSS_IDLE1 && sch) || - (pss->send_state != WSS_SEND_ARTIFACT_INFO && sch && !sch->walk) || - (pss->send_state == WSS_SEND_ARTIFACT_INFO && (!sch || !sch->owner.head))) { - - /* does he want to subscribe to logs? */ - if (sch && sch->logsub && sch->one_task && !pss->subs_list.owner) { - strcpy(pss->sub_task_uuid, sch->one_task->uuid); - lws_dll2_add_head(&pss->subs_list, &pss->vhd->subs_owner); - pss->sub_timestamp = pss->initial_log_timestamp; /* where we got up to */ - lws_callback_on_writable(pss->wsi); - - lwsl_info("%s: subscribed to logs for %s\n", __func__, - pss->sub_task_uuid); + if (lws_ptr_diff_size_t(end, p) < 256) { + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, fi, 0)); + fi = 0; + p = start; } - - pss->send_state = WSS_IDLE1; - saiw_dealloc_sched(sch); } - if (lws_write(pss->wsi, start, lws_ptr_diff_size_t(p, start), - (enum lws_write_protocol)flags) < 0) - return -1; - - lws_callback_on_writable(pss->wsi); - - return 0; - -no_sch: - pss->send_state = WSS_IDLE1; + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, + lws_ptr_diff_size_t(p, start), + lws_write_ws_flags(LWS_WRITE_TEXT, fi, 1)); return 0; } diff --git a/src/web/w-ws-server.c b/src/web/w-ws-server.c index 8bf2bae..6047109 100644 --- a/src/web/w-ws-server.c +++ b/src/web/w-ws-server.c @@ -66,30 +66,6 @@ enum { }; /* - * This allows other parts of sai-web to queue a raw buffer to be sent to - * all connected browsers, eg, for load reports. - * - * The flags are lws_write() flags. - */ -void -saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, - enum lws_write_protocol flags) -{ - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { - struct pss *pss = lws_container_of(p, struct pss, same); - int *pi = (int *)((const char *)buf - sizeof(int)); - - *pi = (int)flags; - - if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) < 0) - lwsl_wsi_err(pss->wsi, "unable to buflist_append"); /* still ask to drain */ - - lws_callback_on_writable(pss->wsi); - - } lws_end_foreach_dll(p); -} - -/* * sai-web is receiving from sai-server * * This may come in chunks and is statefully parsed @@ -139,17 +115,17 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) */ switch (m->a.top_schema_index) { case SAIS_WS_WEBSRV_RX_LOADREPORT: - saiw_ws_broadcast_raw(vhd, buf, len, + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, 0)); break; case SAIS_WS_WEBSRV_RX_TASKACTIVITY: - saiw_ws_broadcast_raw(vhd, buf, len, + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, 0)); break; case SAIS_WS_WEBSRV_RX_SAI_BUILDERS: - saiw_ws_broadcast_raw(vhd, buf, len, + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); @@ -165,7 +141,7 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) case SAIS_WS_WEBSRV_RX_TASKCHANGE: case SAIS_WS_WEBSRV_RX_EVENTCHANGE: case SAIS_WS_WEBSRV_RX_SAI_BUILDERS: - saiw_ws_broadcast_raw(vhd, buf, len, + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); @@ -216,7 +192,7 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { struct pss *pss = lws_container_of(p, struct pss, same); - saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); + saiw_browser_broadcast_queue_builders(pss->vhd, pss); } lws_end_foreach_dll(p); break; @@ -225,7 +201,7 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { struct pss *pss = lws_container_of(p, struct pss, same); - saiw_alloc_sched(pss, WSS_PREPARE_OVERVIEW); + saiw_browser_queue_overview(pss->vhd, pss); } lws_end_foreach_dll(p); break; @@ -234,18 +210,18 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) lws_start_foreach_dll(struct lws_dll2 *, p, vhd->subs_owner.head) { struct pss *pss = lws_container_of(p, struct pss, subs_list); if (!strcmp(pss->sub_task_uuid, ei->event_hash)) - lws_callback_on_writable(pss->wsi); + saiw_broadcast_logs_batch(vhd, pss); } lws_end_foreach_dll(p); break; case SAIS_WS_WEBSRV_RX_LOADREPORT: // lwsl_notice("%s: ^^^^^^^^^^^^^^ SAIS_WS_WEBSRV_RX_LOADREPORT forwarding to browser\n", __func__); // lwsl_hexdump_notice(buf, len); - saiw_ws_broadcast_raw(vhd, buf, len - (unsigned int)n, + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); break; case SAIS_WS_WEBSRV_RX_TASKACTIVITY: - saiw_ws_broadcast_raw(vhd, buf, len - (unsigned int)n, + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); break; }
Page fetched 0s ago, creation time: 9ms (vhost etag hits: 0%, cache hits: 0%)