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 / linux-ubuntu-2004.svg
Author[]Andy Green <andy@warmcat.com> 2025-11-09 12:57 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-11-12 04:45 UTC
Tree46b88c29ffa6d4900d27363d43d30fb9e8556436   Raw Patch
 
power: use modern queue to server
power: use modern queue to server
diff --git a/src/builder/b-ws-server.c b/src/builder/b-ws-server.c index 5de2db6..2f1e344 100644 --- a/src/builder/b-ws-server.c +++ b/src/builder/b-ws-server.c @@ -121,7 +121,7 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, return -1; } - sai_dump_stderr((const char *)buf + LWS_PRE, w); + sai_dump_stderr(buf + LWS_PRE, w); if (saib_srv_queue_tx(h, buf + LWS_PRE, w, ssf)) return -1; diff --git a/src/common/c-utils.c b/src/common/c-utils.c index 671c4e5..f3d85f9 100644 --- a/src/common/c-utils.c +++ b/src/common/c-utils.c @@ -21,6 +21,8 @@ #include <libwebsockets.h> +#include <assert.h> + #include "include/private.h" #if defined(WIN32) @@ -50,7 +52,8 @@ sai_metrics_hash(uint8_t *key, size_t key_len, const char *sp_name, struct lws_genhash_ctx ctx; uint8_t hash[32]; - lwsl_notice("%s: }}}}}}}}}}}}}}}}}}}}} '%s' '%s' '%s' '%s'\n", __func__, sp_name, spawn, project_name, ref); +// lwsl_notice("%s: }}}}}}}}}}}}}}}}}}}}} '%s' '%s' '%s' '%s'\n", __func__, +// sp_name, spawn, project_name, ref); if (lws_genhash_init(&ctx, LWS_GENHASH_TYPE_SHA256) || lws_genhash_update(&ctx, sp_name, strlen(sp_name)) || @@ -88,10 +91,333 @@ sai_task_describe(sai_task_t *task, char *buf, size_t len) } void -sai_dump_stderr(const char *buf, size_t w) +sai_dump_stderr(const uint8_t *buf, size_t w) { if ((ssize_t)write(2, "\n", 1) != (ssize_t)1 || (ssize_t)write(2, buf, LWS_POSIX_LENGTH_CAST(w)) != (ssize_t)w || (ssize_t)write(2, "\n", 1) != (ssize_t)1) lwsl_err("%s: failed to log to stderr\n", __func__); } + + +int +sai_ss_queue_frag_on_buflist_REQUIRES_LWS_PRE(struct lws_ss_handle *h, + struct lws_buflist **buflist, + void *buf, size_t len, + unsigned int ss_flags) +{ + unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); + + *pi = ss_flags; + + if (lws_buflist_append_segment(buflist, (uint8_t *)buf - sizeof(int), + len + sizeof(int)) < 0) + lwsl_ss_err(h, "failed to append"); /* still ask to drain */ + + if (lws_ss_request_tx(h)) + lwsl_ss_err(h, "failed to request tx"); + + return 0; +} + +int +sai_ss_serialize_queue_helper(struct lws_ss_handle *h, + struct lws_buflist **buflist, + const lws_struct_map_t *map, + size_t map_len, void *root) +{ + lws_struct_json_serialize_result_t r = 0; + uint8_t buf[1100 + LWS_PRE], fi = 1; + lws_struct_serialize_t *js; + + js = lws_struct_json_serialize_create(map, map_len, 0, root); + if (!js) { + lwsl_ss_warn(h, "Failed to serialize state update"); + return 1; + } + + do { + size_t w; + + r = lws_struct_json_serialize(js, buf + LWS_PRE, + 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)); + fi = 0; + } while (r == LSJS_RESULT_CONTINUE); + + lws_struct_json_serialize_destroy(&js); + + return 0; +} + +lws_ss_state_return_t +sai_ss_tx_from_buflist_helper(struct lws_ss_handle *ss, struct lws_buflist **buflist, + uint8_t *buf, size_t *len, int *flags) +{ + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(buflist), depi, fl; + char som, som1, eom, final = 1; + size_t fsl, used; + + if (!*buflist) + return LWSSSSRET_TX_DONT_SEND; + + depi = *pi; + + fsl = lws_buflist_next_segment_len(buflist, NULL); + + lws_buflist_fragment_use(buflist, NULL, 0, &som, &eom); + if (som) { + fsl -= sizeof(int); + lws_buflist_fragment_use(buflist, buf, sizeof(int), &som1, &eom); + } + if (!(depi & LWSSS_FLAG_SOM)) + som = 0; + + used = (size_t)lws_buflist_fragment_use(buflist, (uint8_t *)buf, *len, &som1, &eom); + if (!used) + return LWSSSSRET_TX_DONT_SEND; + + if (used < fsl || !(depi & LWSSS_FLAG_EOM)) /* we saved SS flags at the start of the buf */ + final = 0; + + *len = used; + fl = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); + + if ((fl & LWSSS_FLAG_SOM) && (((*flags) & 3) == 2)) { + lwsl_ss_err(ss, "TX: Illegal LWSSS_FLAG_SOM after previous frame without LWSSS_FLAG_EOM"); + assert(0); + } + if (!(fl & LWSSS_FLAG_SOM) && ((*flags) & 3) == 3) { + lwsl_ss_err(ss, "TX: Missing LWSSS_FLAG_SOM after previous frame with LWSSS_FLAG_EOM"); + assert(0); + } + if (!(fl & LWSSS_FLAG_SOM) && !((*flags) & 2)) { + lwsl_ss_err(ss, "TX: Missing LWSSS_FLAG_SOM on first frame"); + assert(0); + } + + *flags = fl; + + /* If there are more to send, request another writable callback */ + if (*buflist && lws_ss_request_tx(ss)) + lwsl_ss_warn(ss, "tx request failed"); + + 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/common/include/private.h b/src/common/include/private.h index f108266..274c7a1 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -1,7 +1,7 @@ /* * Sai - ./src/common/include/private.h * - * Copyright (C) 2019 - 2021 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> * * This library is free software; you can redistribute it and/or * modify it under the terms of the GNU Lesser General Public @@ -18,7 +18,7 @@ * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, * MA 02110-1301 USA * - * structs common to builder and server + * structs common across the various different sai daemons and tools */ #if defined(WIN32) @@ -62,6 +62,14 @@ enum { SAISPRF_SIGNALLED = 0x4000, }; +typedef struct sais_sqlite_cache { + lws_dll2_t list; + char uuid[65]; + sqlite3 *pdb; + lws_usec_t idle_since; + int refcount; +} sais_sqlite_cache_t; + /* The top-level load report message from a builder */ typedef struct sai_active_task_info { lws_dll2_t list; @@ -448,6 +456,38 @@ typedef struct sai_uuid_list { * * It's also used as the object on server side that represents a builder / * platform instance and status. + * + * + * Caution: about the naming, there are `builder platform triplets`, which bind + * to the .sai.json description of the type of builder needed for particular + * tasks. These look like, eg: + * + * rocky9/x86_64-amd/gcc + * coverity/x86_64/gcc + * + * and then there are `builder device names` which represent individual physical + * builder devices which can be powered up and down. These look like, eg + * + * l2 + * ubuntu_rpi4 + * + * One builder can offer multiple different platform triplets. In the examples + * above, the builder l2 offers both rocky9/x86_64-amd/gcc and + * coverity/x86_64/gcc on the same box. + * + * When sai-server looks for a matching builder platform triplet needed by a + * task, it's not bothered which device it binds the job to, if more than one + * offer the platform. So you can have multiple builder devices offering + * popular platforms and jobs should be shared between them. At other times + * sai has to disambiguate the device + platform triplet, in these cases it + * looks like: + * + * l2.rocky9/x86_64-amd/gcc + * l2.coverity/x86_64/gcc + * + * For power purposes, only the builder device name is considered, since we can + * only turn the whole device on or off. If any platform offered by the builder + * is in use, the builder device is kept on. */ typedef struct sai_plat { @@ -561,6 +601,19 @@ typedef struct sai_stay { char stay_on; /* 0 = release, 1 = set */ } sai_stay_t; +typedef struct sai_controlled_builder { + lws_dll2_t list; + char name[64]; +} sai_controlled_builder_t; + +/* sai-power -> sai-server, tells it about power controllers */ +typedef struct sai_power_controller { + lws_dll2_t list; + lws_dll2_owner_t controlled_builders_owner; + char name[64]; + char on; +} sai_power_controller_t; + /* sai-power -> sai-server, tells it the builders it can manage */ typedef struct sai_power_managed_builder { lws_dll2_t list; @@ -571,6 +624,7 @@ typedef struct sai_power_managed_builder { typedef struct sai_power_managed_builders { lws_dll2_t list; lws_dll2_owner_t builders; /* sai_power_managed_builder_t */ + lws_dll2_owner_t power_controllers; /* sai_power_controller_t */ } sai_power_managed_builders_t; @@ -580,13 +634,22 @@ typedef struct sai_stay_state_update { char stay_on; } sai_stay_state_update_t; +/* + * Because the definitions of these arrays of map structs are mostly in + * common/struct-metadata.c, we are forced to repeat the length of the struct + * so we can know the length at the usage. + * + * We must take care to also maintain these lengths when the struct definitions + * change length. + */ extern const lws_struct_map_t lsm_stay[2], lsm_schema_stay[1], lsm_power_managed_builder[2], - lsm_power_managed_builders_list[1], + lsm_power_managed_builders_list[2], lsm_schema_power_managed_builders[1], + lsm_power_controller[3], lsm_schema_json_map_task[], lsm_schema_sq3_map_task[], lsm_schema_sq3_map_event[], @@ -654,4 +717,37 @@ const char * sai_get_ref(const char *fullref); void -sai_dump_stderr(const char *buf, size_t w); +sai_dump_stderr(const uint8_t *buf, size_t w); + +int +sai_ss_queue_frag_on_buflist_REQUIRES_LWS_PRE(struct lws_ss_handle *h, + struct lws_buflist **buflist, + void *buf, size_t len, + unsigned int ss_flags); + +int +sai_ss_serialize_queue_helper(struct lws_ss_handle *h, + struct lws_buflist **buflist, + const lws_struct_map_t *map, + size_t map_len, void *root); + +lws_ss_state_return_t +sai_ss_tx_from_buflist_helper(struct lws_ss_handle *ss, struct lws_buflist **buflist, + uint8_t *buf, size_t *len, int *flags); + +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); +void +sai_event_db_close(lws_dll2_owner_t *sqlite3_cache, sqlite3 **ppdb); + +int +sai_event_db_close_all_now(lws_dll2_owner_t *sqlite3_cache); + +int +sai_event_db_delete_database(const char *sqlite3_path_lhs, const char *event_uuid); + +int +sai_sqlite3_statement(sqlite3 *pdb, const char *cmd, const char *desc); + diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index f5e4d38..5f47757 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -19,6 +19,9 @@ * MA 02110-1301 USA * * lws_struct metadata for structs common to builder and server + * + * For arrays, keep extern length in common/include/private.h in sync + * with changes to array lengths! */ #include <libwebsockets.h> @@ -74,7 +77,7 @@ const lws_struct_map_t lsm_schema_sq3_map_build_metric[] = { }; -const lws_struct_map_t lsm_plat[] = { /* !!! keep extern length in common/include/private.h in sync */ +const lws_struct_map_t lsm_plat[] = { LSM_UNSIGNED (sai_plat_t, uid, "uid"), LSM_STRING_PTR (sai_plat_t, name, "name"), LSM_STRING_PTR (sai_plat_t, platform, "platform"), @@ -287,9 +290,9 @@ const lws_struct_map_t lsm_schema_json_map_artifact[] = { }; const lws_struct_map_t lsm_power_state[] = { - LSM_CARRAY(sai_power_state_t, host, "host"), - LSM_SIGNED(sai_power_state_t, powering_up, "powering_up"), - LSM_SIGNED(sai_power_state_t, powering_down, "powering_down"), + LSM_CARRAY(sai_power_state_t, host, "host"), + LSM_SIGNED(sai_power_state_t, powering_up, "powering_up"), + LSM_SIGNED(sai_power_state_t, powering_down, "powering_down"), }; const lws_struct_map_t lsm_schema_sq3_map_artifact[] = { @@ -297,14 +300,25 @@ const lws_struct_map_t lsm_schema_sq3_map_artifact[] = { }; const lws_struct_map_t lsm_stay[] = { - LSM_CARRAY(sai_stay_t, builder_name, "builder_name"), - LSM_UNSIGNED(sai_stay_t, stay_on, "stay_on"), + LSM_CARRAY(sai_stay_t, builder_name, "builder_name"), + LSM_UNSIGNED(sai_stay_t, stay_on, "stay_on"), }; const lws_struct_map_t lsm_schema_stay[] = { LSM_SCHEMA(sai_stay_t, NULL, lsm_stay, "com.warmcat.sai.power.stay"), }; +const lws_struct_map_t lsm_controlled_builder[] = { + LSM_CARRAY(sai_controlled_builder_t, name, "name"), +}; + +const lws_struct_map_t lsm_power_controller[] = { + LSM_CARRAY(sai_power_controller_t, name, "name"), + LSM_UNSIGNED(sai_power_controller_t, on, "on"), + LSM_LIST(sai_power_controller_t, controlled_builders_owner, + sai_controlled_builder_t, list, NULL, + lsm_controlled_builder, "controlled_builders"), +}; const lws_struct_map_t lsm_power_managed_builder[] = { LSM_CARRAY(sai_power_managed_builder_t, name, "name"), @@ -315,6 +329,9 @@ const lws_struct_map_t lsm_power_managed_builders_list[] = { LSM_LIST(sai_power_managed_builders_t, builders, sai_power_managed_builder_t, list, NULL, lsm_power_managed_builder, "builders"), + LSM_LIST(sai_power_managed_builders_t, power_controllers, + sai_power_controller_t, list, NULL, + lsm_power_controller, "power_controllers"), }; const lws_struct_map_t lsm_schema_power_managed_builders[] = { diff --git a/src/power/CMakeLists.txt b/src/power/CMakeLists.txt index 5942b29..ac91e4a 100644 --- a/src/power/CMakeLists.txt +++ b/src/power/CMakeLists.txt @@ -4,11 +4,12 @@ set(CPACK_DEBIAN_BUILDER_PACKAGE_NAME "sai-power") set(SRCS p-sai.c p-conf.c - p-comms.c p-smartplug.c - p-api.c + p-http-api.c p-utils.c + p-ws-server.c p-tasmota-monitor.c + ../common/c-utils.c ../common/struct-metadata.c ) diff --git a/src/power/p-api.c b/src/power/p-api.c deleted file mode 100644 index 2dbde16..0000000 --- a/src/power/p-api.c +++ /dev/null @@ -1,430 +0,0 @@ -/* - * sai-power - * - * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - * - * This is the h1 API that can be used on the LAN side - */ - -#include <libwebsockets.h> -#include <string.h> -#include <signal.h> -#include <stdlib.h> -#include <sys/stat.h> -#include <fcntl.h> - -#if defined(__linux__) -#include <unistd.h> -#endif - -#if defined(__APPLE__) -#include <sys/stat.h> /* for mkdir() */ -#include <unistd.h> /* for chown() */ -#endif - -#include "p-private.h" - -extern struct lws_spawn_piped *lsp_wol; - -extern struct sai_power power; - - -static void -saip_sul_action_power_off(struct lws_sorted_usec_list *sul) -{ - saip_server_plat_t *sp = lws_container_of(sul, - saip_server_plat_t, sul_delay_off); - lws_ss_state_return_t r; - saip_pcon_t *pc; - - if (!sp->pcon_list.owner) { - lwsl_notice("%s: no power-controller ss for %s\n", __func__, sp->host); - return; - } - - pc = lws_container_of(sp->pcon_list.owner, saip_pcon_t, - controlled_plats_owner); - - saip_notify_server_power_state(sp->host, 0, 1); - - lwsl_warn("%s: powering OFF host %s via power-control %s\n", __func__, sp->host, pc->name); - - r = lws_ss_client_connect(pc->ss_tasmota_off); - if (r) - lwsl_ss_err(pc->ss_tasmota_off, "failed to connect tasmota OFF secure stream: %d", r); -} - -saip_server_plat_t * -find_platform(struct sai_power *pwr, const char *host) -{ - lws_start_foreach_dll(struct lws_dll2 *, px, pwr->sai_server_owner.head) { - saip_server_t *s = lws_container_of(px, saip_server_t, list); - - lws_start_foreach_dll(struct lws_dll2 *, px1, s->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(px1, saip_server_plat_t, list); - - if (!strcmp(host, sp->host)) - return sp; - - } lws_end_foreach_dll(px1); - } lws_end_foreach_dll(px); - - return NULL; -} - -void -saip_notify_server_stay_state(const char *plat_name, int stay_on) -{ - sai_stay_state_update_t *ssu; - saip_server_t *sps; - - /* Find the first (usually only) configured sai-server connection */ - if (!power.sai_server_owner.head) { - lwsl_warn("%s: No sai-server configured to notify\n", __func__); - return; - } - sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); - if (!sps->ss) { - lwsl_warn("%s: Not connected to sai-server to notify\n", __func__); - return; - } - - /* Allocate and queue the notification message */ - ssu = malloc(sizeof(*ssu)); - if (!ssu) - return; - - memset(ssu, 0, sizeof(*ssu)); - lws_strncpy(ssu->builder_name, plat_name, sizeof(ssu->builder_name)); - ssu->stay_on = (char)stay_on; - - /* The per-connection user object for the server link is a saip_server_link_t */ - { - saip_server_link_t *pss = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); - - lws_dll2_add_tail(&ssu->list, &pss->stay_state_update_owner); - } - - /* Request a writable callback to send the message */ - if (lws_ss_request_tx(sps->ss)) - lwsl_ss_warn(sps->ss, "Unable to request tx"); - - lwsl_notice("%s: Queued notification for %s\n", __func__, plat_name); -} - -void -saip_set_stay(const char *builder_name, int stay_on) -{ - saip_server_plat_t *sp = find_platform(&power, builder_name); - saip_server_link_t *pss; - saip_server_t *sps; - - if (!sp) - return; - - sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); - pss = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); - sp->stay = (char)stay_on; - saip_notify_server_stay_state(builder_name, stay_on | sp->needed); - - if (stay_on | sp->needed) - saip_builder_bringup(sps, sp, pss); - else - /* - * power-off is delayed, so we just set the stay flag... - * but let's cancel any pending power-off - */ - lws_sul_cancel(&sp->sul_delay_off); - - /* Find the first (usually only) configured sai-server connection */ - if (!power.sai_server_owner.head) { - lwsl_warn("%s: No sai-server configured to notify\n", __func__); - return; - } - - saip_queue_stay_info(sps, sp, pss); -} - -/* - * local-side h1 server for builders to connect to - */ - -LWS_SS_USER_TYPEDEF - char payload[200]; - size_t size; - size_t pos; -} local_srv_t; - -static lws_ss_state_return_t -local_srv_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, - int *flags) -{ - local_srv_t *g = (local_srv_t *)userobj; - lws_ss_state_return_t r = LWSSSSRET_OK; - - if (g->size == g->pos) - return LWSSSSRET_TX_DONT_SEND; - - if (*len > g->size - g->pos) - *len = g->size - g->pos; - - if (!g->pos) - *flags |= LWSSS_FLAG_SOM; - - memcpy(buf, g->payload + g->pos, *len); - g->pos += *len; - - if (g->pos != g->size) /* more to do */ - r = lws_ss_request_tx(lws_ss_from_user(g)); - else - *flags |= LWSSS_FLAG_EOM; - - lwsl_ss_info(lws_ss_from_user(g), "TX %zu, flags 0x%x, r %d", *len, - (unsigned int)*flags, (int)r); - - return r; -} - -static lws_ss_state_return_t -local_srv_state(void *userobj, void *sh, lws_ss_constate_t state, - lws_ss_tx_ordinal_t ack) -{ - local_srv_t *g = (local_srv_t *)userobj; - sai_power_managed_builders_t *pmb; - sai_power_managed_builder_t *b; - char *path = NULL, pn[128]; - saip_server_link_t *pss; - saip_server_plat_t *sp; - saip_server_t *sps; - int apo = 0; - size_t len; - - // lwsl_ss_user(lws_ss_from_user(g), "state %s", lws_ss_state_name((int)state)); - - switch ((int)state) { - case LWSSSCS_CREATING: - return lws_ss_request_tx(lws_ss_from_user(g)); - - case LWSSSCS_SERVER_TXN: - - lws_ss_get_metadata(lws_ss_from_user(g), "path", (const void **)&path, &len); - // lwsl_ss_user(lws_ss_from_user(g), "LWSSSCS_SERVER_TXN path '%.*s' (%d)", (int)len, path, (int)len); - - /* - * path is containing a string like "/power-off/b32" - * match the last part to a known platform and find out how - * to power that off - */ - - if (lws_ss_set_metadata(lws_ss_from_user(g), "mime", "text/html", 9)) - return LWSSSSRET_DISCONNECT_ME; - - /* - * A transaction is starting on an accepted connection. Say - * that we're OK with the transaction, prepare the user - * object with the response, and request tx to start sending it. - */ - lws_ss_server_ack(lws_ss_from_user(g), 0); - - g->pos = 0; - - if (len == 1 && path[0] == '/') { - /* print controllable platforms */ - - g->size = 0; - - lws_start_foreach_dll(struct lws_dll2 *, px, power.sai_server_owner.head) { - saip_server_t *s = lws_container_of(px, saip_server_t, list); - - lws_start_foreach_dll(struct lws_dll2 *, px1, s->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(px1, saip_server_plat_t, list); - - if (g->size) - g->payload[g->size++] = ','; - g->size = g->size + (size_t)lws_snprintf(g->payload + g->size, sizeof(g->payload) - g->size - 3, "%s", sp->host); - - } lws_end_foreach_dll(px1); - } lws_end_foreach_dll(px); - - g->payload[g->size] = '\0'; - goto bail; - } - - if (len > 6 && !strncmp(path, "/stay/", 6)) { - lws_strnncpy(pn, &path[6], len - 6, sizeof(pn)); - - sp = find_platform(&power, pn); - - if (sp) - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "%c", '0' + (sp->stay | sp->needed)); - else - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "unknown host %s", pn); - goto bail; - } - - if (len > 10 && !strncmp(path, "/power-on/", 10)) { - lws_strnncpy(pn, &path[10], len - 10, sizeof(pn)); - sp = find_platform(&power, pn); - if (!sp) { - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "Unable to find host %s", pn); - goto bail; - } - if (sp->power_on_mac) { - saip_notify_server_power_state(sp->host, 1, 0); - if (write(lws_spawn_get_fd_stdxxx(lsp_wol, 0), - sp->power_on_mac, strlen(sp->power_on_mac)) != - (ssize_t)strlen(sp->power_on_mac)) - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "Write to resume %s failed %d", pn, errno); - else - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "Resumed %s with stay", pn); - sp->stay = 1; - goto bail; - } - - if (sp->pcon_list.owner) { - saip_pcon_t *pc = lws_container_of(sp->pcon_list.owner, - saip_pcon_t, - controlled_plats_owner); - if (lws_ss_client_connect(pc->ss_tasmota_on)) { - lwsl_ss_err(pc->ss_tasmota_on, "failed to connect tasmota ON secure stream"); - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "power-on ss failed create %s", sp->host); - goto bail; - } - } else { - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "no power-controller entry for %s", pn); - goto bail; - } - - lwsl_warn("%s: powered on host %s\n", __func__, sp->host); - sp->stay = 1; /* so builder can understand it's manual */ - saip_notify_server_power_state(sp->host, 1, 0); - - pmb = malloc(sizeof(*pmb)); - if (!pmb) - return 1; - memset(pmb, 0, sizeof(*pmb)); - - b = malloc(sizeof(*b)); - if (!b) { - free(pmb); - return 1; - } - memset(b, 0, sizeof(*b)); - - lws_strncpy(b->name, sp->host, sizeof(b->name)); - b->stay_on = sp->stay; - - sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); - pss = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); - - lws_dll2_add_tail(&b->list, &pmb->builders); - lws_dll2_add_tail(&pmb->list, &pss->managed_builders_owner); - - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "Manually powered on %s", sp->host); - goto bail; - } - - if (len > 16 && !strncmp(path, "/auto-power-off/", 16)) { - apo = 1; - lws_strnncpy(pn, &path[16], len - 16, sizeof(pn)); - goto power_off; - } - - if (len < 11 || strncmp(path, "/power-off/", 11)) { - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "URL path needs to start with /power-off/"); - goto bail; - } - - lws_strnncpy(pn, &path[11], len - 11, sizeof(pn)); - -power_off: - - /* - * Let's have a look at the platform - */ - - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "Unable to find host %s", pn); - - sp = find_platform(&power, pn); - if (sp) { - - if (apo) { - char needs[128]; - - /* - * Since it's not a manual request, - * we should deny it if any deps still need us - */ - - 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); - - if (sp1->needed) - 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); - goto bail; - } - } - - /* - * OK this is it, schedule it to happen - */ - lws_sul_schedule(lws_ss_cx_from_user(g), 0, - &sp->sul_delay_off, - saip_sul_action_power_off, - 3 * LWS_USEC_PER_SEC); - - lwsl_warn("%s: scheduled powering off host %s\n", - __func__, sp->host); - - g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), - "ACK: Scheduled powering off host %s", sp->host); - - sp->stay = 0; /* reset any manual power up */ - } - -bail: - return lws_ss_request_tx_len(lws_ss_from_user(g), - (unsigned long)g->size); - } - - return LWSSSSRET_OK; -} - - -LWS_SS_INFO("local", local_srv_t) - .tx = local_srv_tx, - .state = local_srv_state, -}; diff --git a/src/power/p-comms.c b/src/power/p-comms.c deleted file mode 100644 index afd8fa0..0000000 --- a/src/power/p-comms.c +++ /dev/null @@ -1,483 +0,0 @@ -/* - * sai-power com-warmcat-sai client protocol implementation - * - * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> - * - * This library is free software; you can redistribute it and/or - * modify it under the terms of the GNU Lesser General Public - * License as published by the Free Software Foundation: - * version 2.1 of the License. - * - * This library is distributed in the hope that it will be useful, - * but WITHOUT ANY WARRANTY; without even the implied warranty of - * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU - * Lesser General Public License for more details. - * - * You should have received a copy of the GNU Lesser General Public - * License along with this library; if not, write to the Free Software - * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, - * MA 02110-1301 USA - * - * This is the part of sai-power that handles communication with sai-server - */ - -#include <libwebsockets.h> -#include <string.h> -#include <signal.h> - -#include "p-private.h" - -/* Map for the "powering up" message we send to the server */ -static const lws_struct_map_t lsm_schema_power_state[] = { - LSM_SCHEMA(sai_power_state_t, NULL, lsm_power_state, - "com.warmcat.sai.powerstate"), -}; - -void -saip_notify_server_power_state(const char *plat_name, int up, int down) -{ - saip_server_t *sps; - sai_power_state_t *ps; - - /* Find the first (usually only) configured sai-server connection */ - if (!power.sai_server_owner.head) { - lwsl_warn("%s: No sai-server configured to notify\n", __func__); - return; - } - sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); - if (!sps->ss) { - lwsl_warn("%s: Not connected to sai-server to notify\n", __func__); - return; - } - - /* Allocate and queue the notification message */ - ps = malloc(sizeof(*ps)); - if (!ps) - return; - - memset(ps, 0, sizeof(*ps)); - lws_strncpy(ps->host, plat_name, sizeof(ps->host)); - ps->powering_up = (char)up; - ps->powering_down = (char)down; - - /* The per-connection user object for the server link is a saip_server_link_t */ - { - saip_server_link_t *pss = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); - lws_dll2_add_tail(&ps->list, &pss->ps_owner); - } - - /* Request a writable callback to send the message */ - if (lws_ss_request_tx(sps->ss)) - lwsl_ss_warn(sps->ss, "Unable to request tx"); - - lwsl_notice("%s: Queued notification for %s\n", __func__, plat_name); -} - -int -saip_queue_stay_info(saip_server_t *sps, saip_server_plat_t *sp, saip_server_link_t *pss) -{ - sai_power_managed_builders_t *pmb; - sai_power_managed_builder_t *b; - - pmb = malloc(sizeof(*pmb)); - if (!pmb) - return 1; - - memset(pmb, 0, sizeof(*pmb)); - - /* queue the update for the builder state */ - - b = malloc(sizeof(*b)); - if (!b) { - free(pmb); - return 1; - } - - memset(b, 0, sizeof(*b)); - lws_strncpy(b->name, sp->host, sizeof(b->name)); - b->stay_on = sp->stay; - - lws_dll2_add_tail(&b->list, &pmb->builders); - lws_dll2_add_tail(&pmb->list, &pss->managed_builders_owner); - - if (lws_ss_request_tx(sps->ss)) - lwsl_ss_warn(sps->ss, "Unable to request tx"); - - return 0; -} - -int -saip_builder_bringup(saip_server_t *sps, saip_server_plat_t *sp, - saip_server_link_t *pss) -{ - saip_notify_server_power_state(sp->name, 1, 0); - - if (sp->power_on_type && !strcmp(sp->power_on_type, "wol")) { - lwsl_notice("%s: triggering WOL\n", __func__); - write(lws_spawn_get_fd_stdxxx(lsp_wol, 0), - sp->power_on_mac, strlen(sp->power_on_mac)); - } - - if (sp->pcon_list.owner) { - saip_pcon_t *pc = lws_container_of(sp->pcon_list.owner, - saip_pcon_t, - controlled_plats_owner); - - lwsl_ss_notice(pc->ss_tasmota_on, "starting tasmota"); - if (lws_ss_client_connect(pc->ss_tasmota_on)) - lwsl_ss_err(pc->ss_tasmota_on, "failed to connect tasmota ON secure stream"); - } - - return saip_queue_stay_info(sps, sp, pss); -} - -static lws_ss_state_return_t -saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) -{ - saip_server_link_t *pss = (saip_server_link_t *)userobj; - saip_server_t *sps = (saip_server_t *)lws_ss_opaque_from_user(pss); - const char *p = (const char *)buf, *end = (const char *)buf + len; - char plat[128], benched[4096]; - size_t n, bp = 0; - lws_struct_args_t a; - struct lejp_ctx ctx; - - lwsl_notice("%s: len %d, flags: %d (saip_server_t %p)\n", __func__, (int)len, flags, (void *)sps); - lwsl_hexdump_notice(buf, len); - - memset(&a, 0, sizeof(a)); - a.map_st[0] = lsm_schema_stay; - a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_stay); - a.ac_block_size = 512; - - lws_struct_json_init_parse(&ctx, NULL, &a); - if (lejp_parse(&ctx, (uint8_t *)buf, (int)len) >= 0 && a.dest) { - sai_stay_t *stay = (sai_stay_t *)a.dest; - - // {"schema":"com.warmcat.sai.power.stay","builder_name":"ubuntu_rpi4","stay_on":1} - - lwsl_warn("%s: received stay %s: %d\n", __func__, stay->builder_name, stay->stay_on); - - saip_set_stay(stay->builder_name, stay->stay_on); - lwsac_free(&a.ac); - return 0; - } - lwsac_free(&a.ac); - - /* starting position is that no server-plat is needed */ - - lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); - sp->needed = 0; - } lws_end_foreach_dll(px); - - fprintf(stderr, "|||||||||||||||||||||||||||||||| Server says needed: '%.*s'\n", (int)len, buf); - - while (p < end) { - n = 0; - while (p < end && *p != ',') - if (n < sizeof(plat) - 1) - plat[n++] = *p++; - - plat[n] = '\0'; - if (p < end && *p == ',') - p++; - - /* - * Does this server list this platform as having startable or ongoing - * tasks? - */ - - lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); - - /* - * How about any dependency listed? - */ - - 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); - - lwsl_notice("%s: setting %s as needed dep\n", __func__, sp1->name); - sp1->needed = 2; - saip_set_stay(sp1->name, sp1->stay); - - } lws_end_foreach_dll(px1); - - /* - * Directly listed as needed? - */ - - if (!strcmp(sp->name, plat)) - sp->needed |= 1; - - } lws_end_foreach_dll(px); - } - - /* - * Cascade dependencies up the platforms - */ - - lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); - - if (sp->needed) { - lws_start_foreach_dll(struct lws_dll2 *, py, sp->dependencies_owner.head) { - saip_server_plat_t *spd = lws_container_of(py, saip_server_plat_t, dependencies_list); - - spd->needed |= 2; - - } lws_end_foreach_dll(py); - } - } lws_end_foreach_dll(px); - - /* - * Bringup any directly needed or needed by dependency builders - */ - - lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); - - if (sp->needed) { - lwsl_notice("%s: Needed builders: %s\n", __func__, sp->name); - - /* - * Server said this platform or at least one dependency - * has pending jobs. sai-power config says this builder - * can do jobs on that platform. Let's make sure it - * is powered on. - */ - - saip_builder_bringup(sps, sp, pss); - - } else { - bp += (size_t)lws_snprintf(&benched[bp], sizeof(benched) - bp - 1, "%s%s", !bp ? "" : ", ", sp->name); - benched[sizeof(benched) - 1] = '\0'; - } - - } lws_end_foreach_dll(px); - - if (bp) - lwsl_notice("%s: Benched builders: %s\n", __func__, benched); - - (void)sps; - - return 0; -} - -static lws_ss_state_return_t -saip_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, - int *flags) -{ - saip_server_link_t *pss = (saip_server_link_t *)userobj; - lws_struct_serialize_t *js; - - if (pss->managed_builders_owner.head) { - sai_power_managed_builders_t *pmb = lws_container_of(pss->managed_builders_owner.head, - sai_power_managed_builders_t, list); - - js = lws_struct_json_serialize_create(lsm_schema_power_managed_builders, - LWS_ARRAY_SIZE(lsm_schema_power_managed_builders), 0, pmb); - if (!js) - lwsl_ss_warn(lws_ss_from_user(pss), "Failed to serialize managed builder"); - else - if (lws_struct_json_serialize(js, buf, *len, len) == LSJS_RESULT_FINISH) - *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - - lws_dll2_remove(&pmb->list); - free(pmb); - goto sendify; - } - - if (pss->stay_state_update_owner.head) { - sai_stay_state_update_t *ssu = lws_container_of(pss->stay_state_update_owner.head, - sai_stay_state_update_t, list); - js = lws_struct_json_serialize_create(lsm_schema_stay_state_update, - LWS_ARRAY_SIZE(lsm_schema_stay_state_update), 0, ssu); - if (!js) - lwsl_ss_warn(lws_ss_from_user(pss), "Failed to serialize state update"); - else - if (lws_struct_json_serialize(js, buf, *len, len) == LSJS_RESULT_FINISH) - *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - - lws_dll2_remove(&ssu->list); - free(ssu); - goto sendify; - } - - if (pss->ps_owner.head) { - /* Dequeue the first pending notification */ - sai_power_state_t *ps = lws_container_of(pss->ps_owner.head, sai_power_state_t, list); - - js = lws_struct_json_serialize_create(lsm_schema_power_state, 1, 0, ps); - if (!js) - lwsl_ss_warn(lws_ss_from_user(pss), "Failed to serialize state update"); - else - if (lws_struct_json_serialize(js, buf, *len, len) == LSJS_RESULT_FINISH) - *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - - lws_dll2_remove(&ps->list); - free(ps); - goto sendify; - } - - return LWSSSSRET_TX_DONT_SEND; - -sendify: - lws_struct_json_serialize_destroy(&js); - - /* If there are more to send, request another writable callback */ - if (pss->ps_owner.head || pss->managed_builders_owner.head || pss->stay_state_update_owner.head) - if (lws_ss_request_tx(lws_ss_from_user(pss))) - lwsl_ss_warn(lws_ss_from_user(pss), "tx request failed"); - - lwsl_hexdump_notice(buf, *len); - - return LWSSSSRET_OK; -} - -static int -cleanup_on_ss_destroy(struct lws_dll2 *d, void *user) -{ - saip_server_link_t *pss = (saip_server_link_t *)user; - saip_server_t *sps = (saip_server_t *)lws_ss_opaque_from_user(pss); - - (void)sps; - - - return 0; -} - -static int -cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user) -{ - return 0; -} - -static lws_ss_state_return_t -saip_m_state(void *userobj, void *sh, lws_ss_constate_t state, - lws_ss_tx_ordinal_t ack) -{ - saip_server_link_t *pss = (saip_server_link_t *)userobj; - saip_server_t *sps = (saip_server_t *)lws_ss_opaque_from_user(pss); - const char *pq; - int n; - - lwsl_info("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name(state), - (unsigned int)ack); - - switch (state) { - - case LWSSSCS_CREATING: - - lwsl_info("%s: binding ss to %p %s\n", __func__, sps, sps->url); - - if (lws_ss_set_metadata(sps->ss, "url", sps->url, strlen(sps->url))) - lwsl_warn("%s: unable to set metadata\n", __func__); - - pq = sps->url; - while (*pq && (pq[0] != '/' || pq[1] != '/')) - pq++; - - if (*pq) { - n = 0; - pq += 2; - while (pq[n] && pq[n] != '/') - n++; - } else { - pq = sps->url; - n = (int)strlen(pq); - } -#if 0 - sps->name = sps->url + strlen(sps->url) + 1; - memcpy((char *)sps->name, pq, (unsigned int)n); - ((char *)sps->name)[n] = '\0'; - - while (strchr(sps->name, '.')) - *strchr(sps->name, '.') = '_'; - while (strchr(sps->name, '/')) - *strchr(sps->name, '/') = '_'; -#endif - break; - - case LWSSSCS_DESTROYING: - - /* - * If the logical SS itself is going down, every platform that - * used us to connect to their server and has nspawns are also - * going down - */ - lws_dll2_foreach_safe(&power.sai_server_owner, sps, - cleanup_on_ss_destroy); - - break; - - case LWSSSCS_CONNECTED: - { - saip_server_link_t *pss = (saip_server_link_t *)userobj; - sai_power_managed_builders_t *pmb = malloc(sizeof(*pmb)); - - lwsl_ss_notice(sps->ss, "@@@@@@@@@@@@@@ sai-power CONNECTED to server"); - - if (!pmb) - return LWSSSSRET_DISCONNECT_ME; - - memset(pmb, 0, sizeof(*pmb)); - - lws_start_foreach_dll(struct lws_dll2 *, p, - sps->sai_plat_owner.head) { - saip_server_plat_t *sp = lws_container_of(p, - saip_server_plat_t, list); - sai_power_managed_builder_t *b = malloc(sizeof(*b)); - - if (!b) - continue; - - memset(b, 0, sizeof(*b)); - lws_strncpy(b->name, sp->host, sizeof(b->name)); - b->stay_on = sp->stay; - - lws_dll2_add_tail(&b->list, &pmb->builders); - } lws_end_foreach_dll(p); - - lws_dll2_add_tail(&pmb->list, &pss->managed_builders_owner); - - return lws_ss_request_tx(sps->ss); - } - case LWSSSCS_DISCONNECTED: - lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, pss->ps_owner.head) { - sai_power_state_t *ps = lws_container_of(d, sai_power_state_t, list); - - lws_dll2_remove(&ps->list); - free(ps); - } lws_end_foreach_dll_safe(d, d1); - - /* - * clean up any ongoing spawns related to this connection - */ - - lwsl_info("%s: DISCONNECTED\n", __func__); - lws_dll2_foreach_safe(&power.sai_server_owner, sps, - cleanup_on_ss_disconnect); - break; - - case LWSSSCS_ALL_RETRIES_FAILED: - lwsl_info("%s: LWSSSCS_ALL_RETRIES_FAILED\n", __func__); - return lws_ss_request_tx(sps->ss); - - case LWSSSCS_QOS_ACK_REMOTE: - lwsl_info("%s: LWSSSCS_QOS_ACK_REMOTE\n", __func__); - break; - - default: - break; - } - - return LWSSSSRET_OK; -} - -LWS_SS_INFO("sai_power", saip_server_link_t) - .rx = saip_m_rx, - .state = saip_m_state, - .tx = saip_m_tx, /* We need a TX handler to send messages */ -}; diff --git a/src/power/p-http-api.c b/src/power/p-http-api.c new file mode 100644 index 0000000..8e18fef --- /dev/null +++ b/src/power/p-http-api.c @@ -0,0 +1,403 @@ +/* + * sai-power + * + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * This is the h1 API that can be used on the LAN side + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <stdlib.h> +#include <sys/stat.h> +#include <fcntl.h> + +#if defined(__linux__) +#include <unistd.h> +#endif + +#if defined(__APPLE__) +#include <sys/stat.h> /* for mkdir() */ +#include <unistd.h> /* for chown() */ +#endif + +#include "p-private.h" + +extern struct lws_spawn_piped *lsp_wol; + +extern struct sai_power power; + + +static void +saip_sul_action_power_off(struct lws_sorted_usec_list *sul) +{ + saip_server_plat_t *sp = lws_container_of(sul, + saip_server_plat_t, sul_delay_off); + lws_ss_state_return_t r; + saip_pcon_t *pc; + + if (!sp->pcon_list.owner) { + lwsl_notice("%s: no power-controller ss for %s\n", __func__, sp->host); + return; + } + + pc = lws_container_of(sp->pcon_list.owner, saip_pcon_t, + controlled_plats_owner); + + saip_notify_server_power_state(sp->host, 0, 1); + + lwsl_warn("%s: powering OFF host %s via power-control %s\n", __func__, sp->host, pc->name); + + r = lws_ss_client_connect(pc->ss_tasmota_off); + if (r) + lwsl_ss_err(pc->ss_tasmota_off, "failed to connect tasmota OFF secure stream: %d", r); +} + +saip_server_plat_t * +find_platform(struct sai_power *pwr, const char *host) +{ + lws_start_foreach_dll(struct lws_dll2 *, px, pwr->sai_server_owner.head) { + saip_server_t *s = lws_container_of(px, saip_server_t, list); + + lws_start_foreach_dll(struct lws_dll2 *, px1, s->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(px1, saip_server_plat_t, list); + + if (!strcmp(host, sp->host)) + return sp; + + } lws_end_foreach_dll(px1); + } lws_end_foreach_dll(px); + + return NULL; +} + +void +saip_notify_server_stay_state(const char *plat_name, int stay_on) +{ + sai_stay_state_update_t ssu; + saip_server_link_t *m; + saip_server_t *sps; + + /* Find the first (usually only) configured sai-server connection */ + if (!power.sai_server_owner.head) { + lwsl_warn("%s: No sai-server configured to notify\n", __func__); + return; + } + sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); + if (!sps->ss) { + lwsl_warn("%s: Not connected to sai-server to notify\n", __func__); + return; + } + + m = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); + + memset(&ssu, 0, sizeof(ssu)); + lws_strncpy(ssu.builder_name, plat_name, sizeof(ssu.builder_name)); + ssu.stay_on = (char)stay_on; + + sai_ss_serialize_queue_helper(sps->ss, &m->bl_pwr_to_srv, + lsm_schema_stay_state_update, + LWS_ARRAY_SIZE(lsm_schema_stay_state_update), + &ssu); +} + +void +saip_set_stay(const char *builder_name, int stay_on) +{ + saip_server_plat_t *sp = find_platform(&power, builder_name); + saip_server_link_t *pss; + saip_server_t *sps; + + if (!sp) + return; + + sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); + pss = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); + sp->stay = (char)stay_on; + saip_notify_server_stay_state(builder_name, stay_on | sp->needed); + + if (stay_on | sp->needed) + saip_builder_bringup(sps, sp, pss); + else + /* + * power-off is delayed, so we just set the stay flag... + * but let's cancel any pending power-off + */ + lws_sul_cancel(&sp->sul_delay_off); + + /* Find the first (usually only) configured sai-server connection */ + if (!power.sai_server_owner.head) { + lwsl_warn("%s: No sai-server configured to notify\n", __func__); + return; + } + + saip_queue_stay_info(sps); +} + +/* + * local-side h1 server for builders to connect to + */ + +LWS_SS_USER_TYPEDEF + char payload[200]; + size_t size; + size_t pos; +} local_srv_t; + +static lws_ss_state_return_t +local_srv_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, + int *flags) +{ + local_srv_t *g = (local_srv_t *)userobj; + lws_ss_state_return_t r = LWSSSSRET_OK; + + if (g->size == g->pos) + return LWSSSSRET_TX_DONT_SEND; + + if (*len > g->size - g->pos) + *len = g->size - g->pos; + + if (!g->pos) + *flags |= LWSSS_FLAG_SOM; + + memcpy(buf, g->payload + g->pos, *len); + g->pos += *len; + + if (g->pos != g->size) /* more to do */ + r = lws_ss_request_tx(lws_ss_from_user(g)); + else + *flags |= LWSSS_FLAG_EOM; + + lwsl_ss_info(lws_ss_from_user(g), "TX %zu, flags 0x%x, r %d", *len, + (unsigned int)*flags, (int)r); + + return r; +} + +static lws_ss_state_return_t +local_srv_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + local_srv_t *g = (local_srv_t *)userobj; + char *path = NULL, pn[128]; + saip_server_plat_t *sp; + saip_server_t *sps; + int apo = 0; + size_t len; + + // lwsl_ss_user(lws_ss_from_user(g), "state %s", lws_ss_state_name((int)state)); + + switch ((int)state) { + case LWSSSCS_CREATING: + return lws_ss_request_tx(lws_ss_from_user(g)); + + case LWSSSCS_SERVER_TXN: + + lws_ss_get_metadata(lws_ss_from_user(g), "path", (const void **)&path, &len); + // lwsl_ss_user(lws_ss_from_user(g), "LWSSSCS_SERVER_TXN path '%.*s' (%d)", (int)len, path, (int)len); + + /* + * path is containing a string like "/power-off/b32" + * match the last part to a known platform and find out how + * to power that off + */ + + if (lws_ss_set_metadata(lws_ss_from_user(g), "mime", "text/html", 9)) + return LWSSSSRET_DISCONNECT_ME; + + /* + * A transaction is starting on an accepted connection. Say + * that we're OK with the transaction, prepare the user + * object with the response, and request tx to start sending it. + */ + lws_ss_server_ack(lws_ss_from_user(g), 0); + + g->pos = 0; + + if (len == 1 && path[0] == '/') { + /* print controllable platforms */ + + g->size = 0; + + lws_start_foreach_dll(struct lws_dll2 *, px, power.sai_server_owner.head) { + saip_server_t *s = lws_container_of(px, saip_server_t, list); + + lws_start_foreach_dll(struct lws_dll2 *, px1, s->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(px1, saip_server_plat_t, list); + + if (g->size) + g->payload[g->size++] = ','; + g->size = g->size + (size_t)lws_snprintf(g->payload + g->size, sizeof(g->payload) - g->size - 3, "%s", sp->host); + + } lws_end_foreach_dll(px1); + } lws_end_foreach_dll(px); + + g->payload[g->size] = '\0'; + goto bail; + } + + if (len > 6 && !strncmp(path, "/stay/", 6)) { + lws_strnncpy(pn, &path[6], len - 6, sizeof(pn)); + + sp = find_platform(&power, pn); + + if (sp) + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "%c", '0' + (sp->stay | sp->needed)); + else + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "unknown host %s", pn); + goto bail; + } + + if (len > 10 && !strncmp(path, "/power-on/", 10)) { + + lws_strnncpy(pn, &path[10], len - 10, sizeof(pn)); + sp = find_platform(&power, pn); + if (!sp) { + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "Unable to find host %s", pn); + goto bail; + } + if (sp->power_on_mac) { + saip_notify_server_power_state(sp->host, 1, 0); + if (write(lws_spawn_get_fd_stdxxx(lsp_wol, 0), + sp->power_on_mac, strlen(sp->power_on_mac)) != + (ssize_t)strlen(sp->power_on_mac)) + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "Write to resume %s failed %d", pn, errno); + else + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "Resumed %s with stay", pn); + sp->stay = 1; + goto bail; + } + + if (sp->pcon_list.owner) { + saip_pcon_t *pc = lws_container_of(sp->pcon_list.owner, + saip_pcon_t, + controlled_plats_owner); + if (lws_ss_client_connect(pc->ss_tasmota_on)) { + lwsl_ss_err(pc->ss_tasmota_on, "failed to connect tasmota ON secure stream"); + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "power-on ss failed create %s", sp->host); + goto bail; + } + } else { + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "no power-controller entry for %s", pn); + goto bail; + } + + lwsl_warn("%s: powered on host %s\n", __func__, sp->host); + + sp->stay = 1; /* so builder can understand it's manual */ + saip_notify_server_power_state(sp->host, 1, 0); + + sps = lws_container_of(power.sai_server_owner.head, + saip_server_t, list); + + saip_queue_stay_info(sps); + + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "Manually powered on %s", sp->host); + goto bail; + } + + if (len > 16 && !strncmp(path, "/auto-power-off/", 16)) { + apo = 1; + lws_strnncpy(pn, &path[16], len - 16, sizeof(pn)); + goto power_off; + } + + if (len < 11 || strncmp(path, "/power-off/", 11)) { + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "URL path needs to start with /power-off/"); + goto bail; + } + + lws_strnncpy(pn, &path[11], len - 11, sizeof(pn)); + +power_off: + + /* + * Let's have a look at the platform + */ + + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "Unable to find host %s", pn); + + sp = find_platform(&power, pn); + if (sp) { + + if (apo) { + char needs[128]; + + /* + * Since it's not a manual request, + * we should deny it if any deps still need us + */ + + 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); + + if (sp1->needed) + 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); + goto bail; + } + } + + /* + * OK this is it, schedule it to happen + */ + lws_sul_schedule(lws_ss_cx_from_user(g), 0, + &sp->sul_delay_off, + saip_sul_action_power_off, + 3 * LWS_USEC_PER_SEC); + + lwsl_warn("%s: scheduled powering off host %s\n", + __func__, sp->host); + + g->size = (size_t)lws_snprintf(g->payload, sizeof(g->payload), + "ACK: Scheduled powering off host %s", sp->host); + + sp->stay = 0; /* reset any manual power up */ + } + +bail: + return lws_ss_request_tx_len(lws_ss_from_user(g), + (unsigned long)g->size); + } + + return LWSSSSRET_OK; +} + + +LWS_SS_INFO("local", local_srv_t) + .tx = local_srv_tx, + .state = local_srv_state, +}; diff --git a/src/power/p-private.h b/src/power/p-private.h index 2388ddc..9911d15 100644 --- a/src/power/p-private.h +++ b/src/power/p-private.h @@ -72,7 +72,7 @@ typedef struct tasmota_parse { } tasmota_parse_t; typedef struct saip_pcon { - struct lws_dll2 list; + struct lws_dll2 list; /* sai_power.sai_pcon_owner */ lws_dll2_owner_t controlled_plats_owner; /* saip_server_plat_t */ @@ -87,6 +87,8 @@ typedef struct saip_pcon { struct lws_ss_handle *ss_tasmota_on; struct lws_ss_handle *ss_tasmota_off; struct lws_ss_handle *ss_tasmota_monitor; + + char on; } saip_pcon_t; struct saip_ws_pss; @@ -174,16 +176,14 @@ LWS_SS_USER_TYPEDEF size_t size; size_t pos; - lws_dll2_owner_t ps_owner; - lws_dll2_owner_t managed_builders_owner; - lws_dll2_owner_t stay_state_update_owner; + struct lws_buflist *bl_pwr_to_srv; } saip_server_link_t; extern struct sai_power power; extern const lws_ss_info_t ssi_saip_server_link_t, ssi_saip_smartplug_t; -extern const struct lws_protocols protocol_com_warmcat_sai, protocol_ws_power; +extern const struct lws_protocols protocol_com_warmcat_sai; extern struct lws_spawn_piped *lsp_wol; int saip_config_global(struct sai_power *power, const char *d); @@ -195,13 +195,14 @@ saip_notify_server_power_state(const char *plat_name, int up, int down); void saip_set_stay(const char *builder_name, int stay_on); int -saip_queue_stay_info(saip_server_t *sps, saip_server_plat_t *sp, - saip_server_link_t *pss); +saip_queue_stay_info(saip_server_t *sps); saip_pcon_t * saip_pcon_by_name(struct sai_power *power, const char *name); int -parse_tasmota_status(tasmota_parse_t *tp); +saip_parse_tasmota_status(tasmota_parse_t *tp); int saip_builder_bringup(saip_server_t *sps, saip_server_plat_t *sp, saip_server_link_t *pss); +void +saip_ss_create_tasmota(void); diff --git a/src/power/p-sai.c b/src/power/p-sai.c index cab332a..d5120bc 100644 --- a/src/power/p-sai.c +++ b/src/power/p-sai.c @@ -180,7 +180,6 @@ static const struct lws_protocols protocol_std = { "protocol_std", callback_std, 0, 0 }; static const struct lws_protocols *pprotocols[] = { -// &protocol_ws_power, &protocol_std, NULL }; @@ -404,39 +403,7 @@ int main(int argc, const char **argv) goto bail; } - /* let's create any needed tasmota ss */ - - lws_start_foreach_dll(struct lws_dll2 *, px, power.sai_pcon_owner.head) { - saip_pcon_t *pc = lws_container_of(px, saip_pcon_t, list); - - if (!strcmp(pc->type, "tasmota") && pc->url) { - lws_snprintf(pc->url_on, sizeof(pc->url_on), - "%s/cm?cmnd=Power%%20On", pc->url); - if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, - (void *)pc->url_on, - &pc->ss_tasmota_on, NULL, NULL)) - lwsl_err("%s: %s: failed to create ON smartplug secure stream %s\n", - __func__, pc->name, pc->url_on); - - lws_snprintf(pc->url_off, sizeof(pc->url_off), - "%s/cm?cmnd=Power%%20Off", pc->url); - if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, - (void *)pc->url_off, - &pc->ss_tasmota_off, NULL, NULL)) - lwsl_err("%s: %s: failed to create OFF smartplug secure stream %s\n", - __func__, pc->name, pc->url_off); - - lws_snprintf(pc->url_monitor, sizeof(pc->url_monitor), - "%s?m=1", pc->url); - if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, - (void *)pc->url_monitor, - &pc->ss_tasmota_monitor, NULL, NULL)) - lwsl_err("%s: %s: failed to create MONITOR smartplug secure stream %s\n", - __func__, pc->name, pc->url_monitor); - } - - } lws_end_foreach_dll(px); - + saip_ss_create_tasmota(); { struct lws_spawn_piped_info info; diff --git a/src/power/p-smartplug.c b/src/power/p-smartplug.c index 7f2ca70..af6faf5 100644 --- a/src/power/p-smartplug.c +++ b/src/power/p-smartplug.c @@ -17,6 +17,9 @@ * License along with this library; if not, write to the Free Software * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, * MA 02110-1301 USA + * + * This is the SS link used to communicate with smartplugs (currently + * tasmota only supported) */ #include <libwebsockets.h> @@ -58,3 +61,40 @@ saip_spc_state(void *userobj, void *sh, lws_ss_constate_t state, LWS_SS_INFO("sai_power_smartplug", saip_smartplug_t) .state = saip_spc_state, }; + +void +saip_ss_create_tasmota() +{ + /* let's create any needed tasmota ss */ + + lws_start_foreach_dll(struct lws_dll2 *, px, power.sai_pcon_owner.head) { + saip_pcon_t *pc = lws_container_of(px, saip_pcon_t, list); + + if (!strcmp(pc->type, "tasmota") && pc->url) { + lws_snprintf(pc->url_on, sizeof(pc->url_on), + "%s/cm?cmnd=Power%%20On", pc->url); + if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, + (void *)pc->url_on, + &pc->ss_tasmota_on, NULL, NULL)) + lwsl_err("%s: %s: failed to create ON smartplug secure stream %s\n", + __func__, pc->name, pc->url_on); + + lws_snprintf(pc->url_off, sizeof(pc->url_off), + "%s/cm?cmnd=Power%%20Off", pc->url); + if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, + (void *)pc->url_off, + &pc->ss_tasmota_off, NULL, NULL)) + lwsl_err("%s: %s: failed to create OFF smartplug secure stream %s\n", + __func__, pc->name, pc->url_off); + + lws_snprintf(pc->url_monitor, sizeof(pc->url_monitor), + "%s?m=1", pc->url); + if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, + (void *)pc->url_monitor, + &pc->ss_tasmota_monitor, NULL, NULL)) + lwsl_err("%s: %s: failed to create MONITOR smartplug secure stream %s\n", + __func__, pc->name, pc->url_monitor); + } + + } lws_end_foreach_dll(px); +} diff --git a/src/power/p-tasmota-monitor.c b/src/power/p-tasmota-monitor.c index a6ebcad..8d81e1b 100644 --- a/src/power/p-tasmota-monitor.c +++ b/src/power/p-tasmota-monitor.c @@ -54,7 +54,7 @@ enum { }; int -parse_tasmota_status(tasmota_parse_t *tp) +saip_parse_tasmota_status(tasmota_parse_t *tp) { lws_tokenize_elem e; unsigned int *i; @@ -110,11 +110,14 @@ parse_tasmota_status(tasmota_parse_t *tp) case LWS_TOKZE_INTEGER: if ((tp->match & 0xff) == TOKORD_VOLTAGE) tp->td.voltage_v = (unsigned int)atoi(tp->ts.token); - if ((tp->match >> 8) == TOKORD_ACTIVE && (tp->match & 0xff) == TOKORD_POWER) + if ((tp->match >> 8) == TOKORD_ACTIVE && + (tp->match & 0xff) == TOKORD_POWER) tp->td.active_power_w = (unsigned int)atoi(tp->ts.token); - if ((tp->match >> 8) == TOKORD_APPARENT && (tp->match & 0xff) == TOKORD_POWER) + if ((tp->match >> 8) == TOKORD_APPARENT && + (tp->match & 0xff) == TOKORD_POWER) tp->td.apparent_power_va = (unsigned int)atoi(tp->ts.token); - if ((tp->match >> 8) == TOKORD_REACTIVE && (tp->match & 0xff) == TOKORD_POWER) + if ((tp->match >> 8) == TOKORD_REACTIVE && + (tp->match & 0xff) == TOKORD_POWER) tp->td.reactive_power_var = (unsigned int)atoi(tp->ts.token); break; @@ -123,13 +126,17 @@ parse_tasmota_status(tasmota_parse_t *tp) if ((tp->match & 0xff) == TOKORD_CURRENT) i = &tp->td.current_ma; - if ((tp->match >> 8) == TOKORD_POWER && (tp->match & 0xff) == TOKORD_FACTOR) + if ((tp->match >> 8) == TOKORD_POWER && + (tp->match & 0xff) == TOKORD_FACTOR) i = &tp->td.power_factor_scaled_1000; - if ((tp->match >> 8) == TOKORD_ENERGY && (tp->match & 0xff) == TOKORD_TODAY) + if ((tp->match >> 8) == TOKORD_ENERGY && + (tp->match & 0xff) == TOKORD_TODAY) i = &tp->td.energy_today_wh; - if ((tp->match >> 8) == TOKORD_ENERGY && (tp->match & 0xff) == TOKORD_YESTERDAY) + if ((tp->match >> 8) == TOKORD_ENERGY && + (tp->match & 0xff) == TOKORD_YESTERDAY) i = &tp->td.energy_yesterday_wh; - if ((tp->match >> 8) == TOKORD_ENERGY && (tp->match & 0xff) == TOKORD_TOTAL) + if ((tp->match >> 8) == TOKORD_ENERGY && + (tp->match & 0xff) == TOKORD_TOTAL) i = &tp->td.energy_total_wh; if (i) { diff --git a/src/power/p-ws-server.c b/src/power/p-ws-server.c new file mode 100644 index 0000000..016240f --- /dev/null +++ b/src/power/p-ws-server.c @@ -0,0 +1,406 @@ +/* + * sai-power com-warmcat-sai client protocol implementation + * + * Copyright (C) 2019 - 2025 Andy Green <andy@warmcat.com> + * + * This library is free software; you can redistribute it and/or + * modify it under the terms of the GNU Lesser General Public + * License as published by the Free Software Foundation: + * version 2.1 of the License. + * + * This library is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU + * Lesser General Public License for more details. + * + * You should have received a copy of the GNU Lesser General Public + * License along with this library; if not, write to the Free Software + * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, + * MA 02110-1301 USA + * + * This is the part of sai-power that handles communication with sai-server + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <assert.h> + +#include "p-private.h" + +/* Map for the "powering up" message we send to the server */ +static const lws_struct_map_t lsm_schema_power_state[] = { + LSM_SCHEMA(sai_power_state_t, NULL, lsm_power_state, + "com.warmcat.sai.powerstate"), +}; + +void +saip_notify_server_power_state(const char *plat_name, int up, int down) +{ + saip_server_link_t *m; + sai_power_state_t ps; + saip_server_t *sps; + + /* Find the first (usually only) configured sai-server connection */ + if (!power.sai_server_owner.head) { + lwsl_warn("%s: No sai-server configured to notify\n", __func__); + return; + } + sps = lws_container_of(power.sai_server_owner.head, saip_server_t, list); + if (!sps->ss) { + lwsl_warn("%s: Not connected to sai-server to notify\n", __func__); + return; + } + + m = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); + + memset(&ps, 0, sizeof(ps)); + + lws_strncpy(ps.host, plat_name, sizeof(ps.host)); + ps.powering_up = (char)up; + ps.powering_down = (char)down; + + sai_ss_serialize_queue_helper(sps->ss, &m->bl_pwr_to_srv, + lsm_schema_power_state, + LWS_ARRAY_SIZE(lsm_schema_power_state), + &ps); +} + +int +saip_queue_stay_info(saip_server_t *sps) +{ + sai_power_managed_builders_t pmb; + struct lwsac *ac = NULL; + saip_server_link_t *m; + int r; + + lwsl_ss_notice(sps->ss, "@@@@@@@@@@@@@@ sai-power CONNECTED to server"); + + m = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); + + memset(&pmb, 0, sizeof(pmb)); + + lws_start_foreach_dll(struct lws_dll2 *, p, sps->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(p, + saip_server_plat_t, list); + sai_power_managed_builder_t *b = lwsac_use_zero(&ac, sizeof(*b), 2048); + + if (b) { + lws_strncpy(b->name, sp->host, sizeof(b->name)); + b->stay_on = sp->stay; + + lws_dll2_add_tail(&b->list, &pmb.builders); + } + } lws_end_foreach_dll(p); + + lws_start_foreach_dll(struct lws_dll2 *, p, power.sai_pcon_owner.head) { + saip_pcon_t *pc = lws_container_of(p, saip_pcon_t, list); + sai_power_controller_t *pc1 = lwsac_use_zero(&ac, sizeof(*pc1), 2048); + + if (pc1) { + lws_strncpy(pc1->name, pc->name, sizeof(pc1->name)); + pc1->on = pc->on; + + lws_dll2_add_tail(&pc1->list, &pmb.power_controllers); + + lws_start_foreach_dll(struct lws_dll2 *, p1, + pc->controlled_plats_owner.head) { + saip_server_plat_t *sp = lws_container_of(p1, + saip_server_plat_t, pcon_list); + sai_controlled_builder_t *c = lwsac_use_zero(&ac, + sizeof(*c), 2048); + + if (c) { + if (sp->host) + lws_strncpy(c->name, sp->host, + sizeof(c->name)); + + lws_dll2_add_tail(&c->list, + &pc1->controlled_builders_owner); + } + } lws_end_foreach_dll(p1); + } + } lws_end_foreach_dll(p); + + r = sai_ss_serialize_queue_helper(sps->ss, &m->bl_pwr_to_srv, + lsm_schema_power_managed_builders, + LWS_ARRAY_SIZE(lsm_schema_power_managed_builders), + &pmb); + lwsac_free(&ac); + + return r; +} + +int +saip_builder_bringup(saip_server_t *sps, saip_server_plat_t *sp, + saip_server_link_t *pss) +{ + saip_notify_server_power_state(sp->name, 1, 0); + + if (sp->power_on_type && !strcmp(sp->power_on_type, "wol")) { + lwsl_notice("%s: triggering WOL\n", __func__); + write(lws_spawn_get_fd_stdxxx(lsp_wol, 0), + sp->power_on_mac, strlen(sp->power_on_mac)); + } + + if (sp->pcon_list.owner) { + saip_pcon_t *pc = lws_container_of(sp->pcon_list.owner, + saip_pcon_t, + controlled_plats_owner); + + lwsl_ss_notice(pc->ss_tasmota_on, "starting tasmota"); + if (lws_ss_client_connect(pc->ss_tasmota_on)) + lwsl_ss_err(pc->ss_tasmota_on, "failed to connect tasmota ON secure stream"); + } + + return saip_queue_stay_info(sps); +} + +static lws_ss_state_return_t +saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + saip_server_link_t *pss = (saip_server_link_t *)userobj; + saip_server_t *sps = (saip_server_t *)lws_ss_opaque_from_user(pss); + const char *p = (const char *)buf, *end = (const char *)buf + len; + char plat[128], benched[4096]; + size_t n, bp = 0; + lws_struct_args_t a; + struct lejp_ctx ctx; + + lwsl_notice("%s: len %d, flags: %d (saip_server_t %p)\n", __func__, (int)len, flags, (void *)sps); + lwsl_hexdump_notice(buf, len); + + memset(&a, 0, sizeof(a)); + a.map_st[0] = lsm_schema_stay; + a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_stay); + a.ac_block_size = 512; + + lws_struct_json_init_parse(&ctx, NULL, &a); + if (lejp_parse(&ctx, (uint8_t *)buf, (int)len) >= 0 && a.dest) { + sai_stay_t *stay = (sai_stay_t *)a.dest; + + // {"schema":"com.warmcat.sai.power.stay","builder_name":"ubuntu_rpi4","stay_on":1} + + lwsl_warn("%s: received stay %s: %d\n", __func__, stay->builder_name, stay->stay_on); + + saip_set_stay(stay->builder_name, stay->stay_on); + lwsac_free(&a.ac); + return 0; + } + lwsac_free(&a.ac); + + /* starting position is that no server-plat is needed */ + + lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); + sp->needed = 0; + } lws_end_foreach_dll(px); + + fprintf(stderr, "|||||||||||||||||||||||||||||||| Server says needed: '%.*s'\n", (int)len, buf); + + while (p < end) { + n = 0; + while (p < end && *p != ',') + if (n < sizeof(plat) - 1) + plat[n++] = *p++; + + plat[n] = '\0'; + if (p < end && *p == ',') + p++; + + /* + * Does this server list this platform as having startable or ongoing + * tasks? + */ + + lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); + + /* + * How about any dependency listed? + */ + + 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); + + lwsl_notice("%s: setting %s as needed dep\n", __func__, sp1->name); + sp1->needed = 2; + saip_set_stay(sp1->name, sp1->stay); + + } lws_end_foreach_dll(px1); + + /* + * Directly listed as needed? + */ + + if (!strcmp(sp->name, plat)) + sp->needed |= 1; + + } lws_end_foreach_dll(px); + } + + /* + * Cascade dependencies up the platforms + */ + + lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); + + if (sp->needed) { + lws_start_foreach_dll(struct lws_dll2 *, py, sp->dependencies_owner.head) { + saip_server_plat_t *spd = lws_container_of(py, + saip_server_plat_t, dependencies_list); + + spd->needed |= 2; + + } lws_end_foreach_dll(py); + } + } lws_end_foreach_dll(px); + + /* + * Bringup any directly needed or needed by dependency builders + */ + + lws_start_foreach_dll(struct lws_dll2 *, px, sps->sai_plat_owner.head) { + saip_server_plat_t *sp = lws_container_of(px, saip_server_plat_t, list); + + if (sp->needed) { + lwsl_notice("%s: Needed builders: %s\n", __func__, sp->name); + + /* + * Server said this platform or at least one dependency + * has pending jobs. sai-power config says this builder + * can do jobs on that platform. Let's make sure it + * is powered on. + */ + + saip_builder_bringup(sps, sp, pss); + + } else { + bp += (size_t)lws_snprintf(&benched[bp], sizeof(benched) - bp - 1, + "%s%s", !bp ? "" : ", ", sp->name); + benched[sizeof(benched) - 1] = '\0'; + } + + } lws_end_foreach_dll(px); + + if (bp) + lwsl_notice("%s: Benched builders: %s\n", __func__, benched); + + (void)sps; + + return 0; +} + +static lws_ss_state_return_t +saip_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, + int *flags) +{ + saip_server_link_t *pss = (saip_server_link_t *)userobj; + lws_ss_state_return_t r; + + /* + * helper fills tx with next buflist content, and asks to write again + * if any left. + */ + + r = sai_ss_tx_from_buflist_helper(pss->ss, &pss->bl_pwr_to_srv, + buf, len, flags); + + if (r == LWSSSSRET_OK) + sai_dump_stderr(buf, *len); + + return r; +} + +static int +cleanup_on_ss_destroy(struct lws_dll2 *d, void *user) +{ + saip_server_link_t *pss = (saip_server_link_t *)user; + saip_server_t *sps = (saip_server_t *)lws_ss_opaque_from_user(pss); + + (void)sps; + + + return 0; +} + +static int +cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user) +{ + return 0; +} + +static lws_ss_state_return_t +saip_m_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + saip_server_link_t *pss = (saip_server_link_t *)userobj; + saip_server_t *sps = (saip_server_t *)lws_ss_opaque_from_user(pss); + const char *pq; + int n; + + // lwsl_info("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name(state), + // (unsigned int)ack); + + switch (state) { + + case LWSSSCS_CREATING: + + lwsl_info("%s: binding ss to %p %s\n", __func__, sps, sps->url); + + if (lws_ss_set_metadata(sps->ss, "url", sps->url, strlen(sps->url))) + lwsl_warn("%s: unable to set metadata\n", __func__); + + pq = sps->url; + while (*pq && (pq[0] != '/' || pq[1] != '/')) + pq++; + + if (*pq) { + n = 0; + pq += 2; + while (pq[n] && pq[n] != '/') + n++; + } else { + pq = sps->url; + n = (int)strlen(pq); + } + break; + + case LWSSSCS_DESTROYING: + lws_dll2_foreach_safe(&power.sai_server_owner, sps, + cleanup_on_ss_destroy); + break; + + case LWSSSCS_CONNECTED: + lwsl_ss_notice(sps->ss, "@@@@@@@@@@@@@@ sai-power CONNECTED to server"); + saip_queue_stay_info(sps); + break; + + case LWSSSCS_DISCONNECTED: + lws_buflist_destroy_all_segments(&pss->bl_pwr_to_srv); + lwsl_info("%s: DISCONNECTED\n", __func__); + lws_dll2_foreach_safe(&power.sai_server_owner, sps, + cleanup_on_ss_disconnect); + break; + + case LWSSSCS_ALL_RETRIES_FAILED: + lwsl_info("%s: LWSSSCS_ALL_RETRIES_FAILED\n", __func__); + return lws_ss_request_tx(sps->ss); + + case LWSSSCS_QOS_ACK_REMOTE: + lwsl_info("%s: LWSSSCS_QOS_ACK_REMOTE\n", __func__); + break; + + default: + break; + } + + return LWSSSSRET_OK; +} + +LWS_SS_INFO("sai_power", saip_server_link_t) + .rx = saip_m_rx, + .state = saip_m_state, + .tx = saip_m_tx, +}; diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index df2857d..576dfa0 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -20,6 +20,16 @@ set(SRCS ../common/struct-metadata.c ) +include(CheckCCompilerFlag) +include(CheckFunctionExists) +include(CheckSymbolExists) +include(CheckIncludeFile) +include(CheckIncludeFiles) +include(CheckLibraryExists) +include(CheckTypeSize) +include(CheckCSourceCompiles) +include(GNUInstallDirs) + set(requirements 1) require_lws_config(LWS_WITH_SERVER 1 requirements) require_lws_config(LWS_WITH_GENCRYPTO 1 requirements) @@ -49,6 +59,8 @@ if (requirements) include_directories(BEFORE "${SAI_LWS_INC_PATH}") + # set(CMAKE_C_FLAGS "-ffunction-sections -Wl,--gc-sections -Wl,--print-gc-sections ${CMAKE_C_FLAGS}" ) + CHECK_C_SOURCE_COMPILES("#include <libwebsockets.h>\nint main(void) {\n#if defined(LWS_HAVE_LIBCAP)\n return 0;\n#else\n fail;\n#endif\n return 0;\n}\n" HAS_LIBCAP) diff --git a/src/server/s-central.c b/src/server/s-central.c index 924960b..15a499e 100644 --- a/src/server/s-central.c +++ b/src/server/s-central.c @@ -90,7 +90,8 @@ sais_central_clean_abandoned(struct vhd *vhd) sai_event_t *e = lws_container_of(p, sai_event_t, list); sqlite3 *pdb = NULL; - if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { + if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) { char *err = NULL; /* @@ -138,7 +139,7 @@ sais_central_clean_abandoned(struct vhd *vhd) sqlite3_finalize(sm); } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); } } lws_end_foreach_dll(p); diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 59a36ca..4dbceae 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -464,7 +464,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, } if (pss->pdb_artifact) { - sais_event_db_close(pss->vhd, &pss->pdb_artifact); + sai_event_db_close(&pss->vhd->sqlite3_cache, &pss->pdb_artifact); pss->pdb_artifact = NULL; } @@ -520,11 +520,8 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, break; } - if (pss->is_power || pss->stay_owner.head) { - lwsl_notice("%s: going down power tx path\n", __func__); - + if (pss->is_power || pss->stay_owner.head) return sais_power_tx(vhd, pss, buf, sizeof(buf)); - } return sais_ws_json_tx_builder(vhd, pss, buf, sizeof(buf)); diff --git a/src/server/s-helpers.c b/src/server/s-helpers.c index a605bab..ee76cfc 100644 --- a/src/server/s-helpers.c +++ b/src/server/s-helpers.c @@ -71,216 +71,6 @@ sai_task_uuid_to_event_uuid(char *event_uuid33, const char *task_uuid65) 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 @@ -329,7 +119,7 @@ sais_server_destroy(struct vhd *vhd, sais_t *server) lws_dll2_foreach_safe(&server->builder_owner, NULL, sai_detach_builder); - sais_event_db_close_all_now(vhd); + sai_event_db_close_all_now(&vhd->sqlite3_cache); lws_struct_sq3_close(&server->pdb); diff --git a/src/server/s-notification.c b/src/server/s-notification.c index cb0358f..45a5ab9 100644 --- a/src/server/s-notification.c +++ b/src/server/s-notification.c @@ -427,7 +427,8 @@ next_plat: ; * configuration's tasks for each platform */ - if (sais_event_db_ensure_open(pss->vhd, pss->sn.e.uuid, 1, &pdb)) { + if (sai_event_db_ensure_open(pss->vhd->context, &pss->vhd->sqlite3_cache, + pss->vhd->sqlite3_path_lhs, pss->sn.e.uuid, 1, &pdb)) { lwsl_err("%s: unable to open event-specific db\n", __func__); return -1; } @@ -537,7 +538,7 @@ next_plat: ; sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); if (err) sqlite3_free(err); - sais_event_db_close(pss->vhd, &pdb); + sai_event_db_close(&pss->vhd->sqlite3_cache, &pdb); return -1; } @@ -602,7 +603,7 @@ next_plat: ; if (err) sqlite3_free(err); - sais_event_db_close(pss->vhd, &pdb); + sai_event_db_close(&pss->vhd->sqlite3_cache, &pdb); /* * Recompute startable task platforms and broadcast to all sai-power, @@ -1009,7 +1010,8 @@ sai_notification_file_upload_cb(void *data, const char *name, */ sai_uuid16_create(lws_get_context(pss->wsi), pss->sn.e.uuid); - m = sais_event_db_ensure_open(pss->vhd, pss->sn.e.uuid, 1, + m = sai_event_db_ensure_open(pss->vhd->context, &pss->vhd->sqlite3_cache, + pss->vhd->sqlite3_path_lhs, pss->sn.e.uuid, 1, (sqlite3 **)&pss->sn.e.pdb); if (m) { lwsl_err("%s: XX %d unable to open event-specific database\n", @@ -1023,7 +1025,7 @@ sai_notification_file_upload_cb(void *data, const char *name, LWS_ARRAY_SIZE(saifile_paths)); m = lejp_parse(&saictx, (uint8_t *)pss->sn.saifile, (int)pss->sn.saifile_out_pos); - sais_event_db_close(pss->vhd, (sqlite3 **)&pss->sn.e.pdb); + sai_event_db_close(&pss->vhd->sqlite3_cache, (sqlite3 **)&pss->sn.e.pdb); if (m < 0) { lwsl_notice("%s: saifile JSON 1 decode failed '%s' (%d)\n", __func__, lejp_error_to_string(m), m); diff --git a/src/server/s-power.c b/src/server/s-power.c index a941e3f..d61d4dd 100644 --- a/src/server/s-power.c +++ b/src/server/s-power.c @@ -191,7 +191,7 @@ sais_power_tx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) lws_struct_json_serialize_destroy(&js); lwsl_wsi_notice(pss->wsi, "%s: server issuing stay notice\n", __func__); - sai_dump_stderr((char *)start, w); + sai_dump_stderr(start, w); lws_dll2_remove(&s->list); free(s); diff --git a/src/server/s-private.h b/src/server/s-private.h index 14f1a88..757e939 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -195,14 +195,6 @@ struct pss { uint8_t ovstate; /* SOS_ substate when doing overview */ }; -typedef struct sais_sqlite_cache { - lws_dll2_t list; - char uuid[65]; - sqlite3 *pdb; - lws_usec_t idle_since; - int refcount; -} sais_sqlite_cache_t; - typedef struct sais_plat { lws_dll2_t list; const char *plat; @@ -256,21 +248,6 @@ sai_notification_file_upload_cb(void *data, const char *name, enum lws_spa_fileupload_states state); int -sai_sqlite3_statement(sqlite3 *pdb, const char *cmd, const char *desc); - -int -sais_event_db_ensure_open(struct vhd *vhd, const char *event_uuid, char can_create, sqlite3 **ppdb); - -void -sais_event_db_close(struct vhd *vhd, sqlite3 **ppdb); - -int -sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid); - -int -sais_event_db_close_all_now(struct vhd *vhd); - -int sai_sq3_event_lookup(sqlite3 *pdb, uint64_t start, lws_struct_args_cb cb, void *ca); int diff --git a/src/server/s-task-helpers.c b/src/server/s-task-helpers.c index e77d78b..d78ab29 100644 --- a/src/server/s-task-helpers.c +++ b/src/server/s-task-helpers.c @@ -108,7 +108,8 @@ sais_bind_task_to_builder(struct vhd *vhd, const char *builder_name, * Open the event-specific database on the temporary event object */ - if (sais_event_db_ensure_open(vhd, event_uuid, 0, (sqlite3 **)&e->pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, (sqlite3 **)&e->pdb)) { lwsl_err("%s: unable to open event-specific database\n", __func__); @@ -144,7 +145,7 @@ sais_bind_task_to_builder(struct vhd *vhd, const char *builder_name, bail: if (e) - sais_event_db_close(vhd, (sqlite3 **)&e->pdb); + sai_event_db_close(&vhd->sqlite3_cache, (sqlite3 **)&e->pdb); lwsac_free(&ac); return r; @@ -190,7 +191,8 @@ sais_set_task_state(struct vhd *vhd, const char *task_uuid, * Open the event-specific database on the temporary event object */ - if (sais_event_db_ensure_open(vhd, event_uuid, 0, (sqlite3 **)&e->pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, (sqlite3 **)&e->pdb)) { lwsl_err("%s: unable to open event-specific database\n", __func__); @@ -342,7 +344,7 @@ sais_set_task_state(struct vhd *vhd, const char *task_uuid, } } - sais_event_db_close(vhd, (sqlite3 **)&e->pdb); + sai_event_db_close(&vhd->sqlite3_cache, (sqlite3 **)&e->pdb); lwsac_free(&ac); if (ostate == SAIES_STEP_SUCCESS) { @@ -354,7 +356,7 @@ sais_set_task_state(struct vhd *vhd, const char *task_uuid, bail: if (e) - sais_event_db_close(vhd, (sqlite3 **)&e->pdb); + sai_event_db_close(&vhd->sqlite3_cache, (sqlite3 **)&e->pdb); lwsac_free(&ac); return 1; @@ -417,7 +419,8 @@ sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) sai_task_uuid_to_event_uuid(event_uuid, task_uuid); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { lwsl_err("%s: unable to open event-specific database\n", __func__); return -1; } @@ -429,14 +432,14 @@ sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) if (sqlite3_exec(pdb, q, sql3_get_string_cb, builder_name, NULL) != SQLITE_OK || !builder_name[0]) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &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); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); /* * This frees the sqlite task from being bound to any builder @@ -494,7 +497,8 @@ sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_ sai_task_uuid_to_event_uuid(event_uuid, task_uuid); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { lwsl_err("%s: unable to open event-specific database\n", __func__); @@ -507,7 +511,7 @@ sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_ ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; lwsl_err("%s: %s: %s: fail\n", __func__, cmd, @@ -519,7 +523,7 @@ sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_ ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; lwsl_err("%s: %s: %s: fail\n", __func__, cmd, @@ -527,7 +531,7 @@ sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_ return SAI_DB_RESULT_ERROR; } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); /* 1,1 == reset started and duration in db for task to 0 */ sais_set_task_state(vhd, task_uuid, SAIES_WAITING, 1, 1); @@ -574,7 +578,8 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) sai_task_uuid_to_event_uuid(event_uuid, task_uuid); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { lwsl_err("%s: unable to open event-specific database\n", __func__); @@ -586,7 +591,7 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) 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); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return SAI_DB_RESULT_ERROR; } @@ -600,7 +605,7 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) ret = sqlite3_exec(pdb, cmd, NULL, NULL, NULL); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; @@ -612,7 +617,7 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid) } lwsac_free(&ac); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); sais_set_task_state(vhd, task_uuid, SAIES_WAITING, 0, 0); diff --git a/src/server/s-task.c b/src/server/s-task.c index 8b158ed..02f5352 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -40,7 +40,8 @@ sais_event_check_for_plat_tasks(struct vhd *vhd, const char *event_uuid, char query[256]; unsigned int count = 0; - if (sais_event_db_ensure_open(vhd, event_uuid, 1, &check_pdb)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 1, &check_pdb)) return 0; lws_snprintf(query, sizeof(query), @@ -51,7 +52,7 @@ sais_event_check_for_plat_tasks(struct vhd *vhd, const char *event_uuid, NULL) != SQLITE_OK) count = 0; - sais_event_db_close(vhd, &check_pdb); + sai_event_db_close(&vhd->sqlite3_cache, &check_pdb); lwsl_notice("%s: event %s, platform %s: count %u\n", __func__, event_uuid, platform, count); @@ -223,7 +224,8 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, // lwsl_notice("candidate event %s '%s'\n", e->uuid, esc_plat); - if (sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) goto next; /* @@ -301,7 +303,8 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, } while (1); if (checked_uuid[0] && - !sais_event_db_ensure_open(vhd, checked_uuid, 1, &prev_pdb)) { + !sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, checked_uuid, 1, &prev_pdb)) { sqlite3_stmt *sm; /* we are looking for failed tasks here */ @@ -347,7 +350,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, } else lwsl_err("%s: query fail 1\n", __func__); - sais_event_db_close(vhd, &prev_pdb); + sai_event_db_close(&vhd->sqlite3_cache, &prev_pdb); } /* @@ -377,7 +380,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, lwsl_notice("%s: Prioritizing failed task for %s ('%s')\n", __func__, platform, fti->taskname); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); lwsac_free(&failed_ac); memcpy(&pss->alloc_task, lws_container_of( @@ -407,7 +410,7 @@ next1: ; goto close_next; lwsl_notice("%s: orig exit\n", __func__); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); lwsac_free(&failed_ac); memcpy(&pss->alloc_task, lws_container_of( @@ -417,7 +420,7 @@ next1: ; return &pss->alloc_task; close_next: - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); next: ; } lws_end_foreach_dll(p); @@ -522,7 +525,8 @@ sais_platforms_with_tasks_pending(struct vhd *vhd) sqlite3_stmt *sm; int n; - if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { + if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) { if (sqlite3_prepare_v2(pdb, "select distinct platform " "from tasks where " @@ -553,7 +557,7 @@ sais_platforms_with_tasks_pending(struct vhd *vhd) __func__, n, sqlite3_errmsg(pdb)); } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); } } lws_end_foreach_dll(p); @@ -680,7 +684,8 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) sai_event_t *e = lws_container_of(d, sai_event_t, list); sqlite3 *pdb = NULL; - if (sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) goto next; if (lws_struct_sq3_deserialize(pdb, @@ -727,7 +732,7 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) next1: lwsac_free(&ac_tasks); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); next: ; } lws_end_foreach_dll(d); @@ -782,7 +787,8 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) event_uuid[0] = '\0'; sai_task_uuid_to_event_uuid(event_uuid, task_uuid); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb) || !pdb) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb) || !pdb) return -1; // lwsl_notice("%s: task_uuid %s, pdb %p\n", __func__, task_uuid, pdb); @@ -793,7 +799,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) lsm_schema_sq3_map_task, &o, &ac, 0, 1); if (n < 0 || !o.head) { lwsl_warn("%s: bailing as nothing with state != 4\n", __func__); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return -1; } @@ -808,7 +814,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) temp_task = malloc(sizeof(sai_task_t)); if (!temp_task) { lwsac_free(&ac); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); return -1; } @@ -828,7 +834,7 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) &temp_task->ac_task_container, 0, 1); if (n < 0 || !o_event.head) { lwsl_warn("%s: bailing as nothing with uuid %s\n", __func__, esc_uuid); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); free(temp_task); return -1; } @@ -954,12 +960,12 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) lws_dll2_add_tail(&temp_task->pending_assign_list, &pss->issue_task_owner); lws_callback_on_writable(pss->wsi); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); return 0; bail: - sais_event_db_close(vhd, &pdb); +sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&temp_task->ac_task_container); free(temp_task); diff --git a/src/server/s-webops.c b/src/server/s-webops.c index 87ba5e5..05e87c2 100644 --- a/src/server/s-webops.c +++ b/src/server/s-webops.c @@ -198,7 +198,8 @@ sais_event_reset(struct vhd *vhd, const char *event_uuid) char *err = NULL; int ret; - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) return SAI_DB_RESULT_ERROR; if (lws_struct_sq3_deserialize(pdb, NULL, NULL, @@ -207,7 +208,7 @@ sais_event_reset(struct vhd *vhd, const char *event_uuid) ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; @@ -219,7 +220,7 @@ sais_event_reset(struct vhd *vhd, const char *event_uuid) 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); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return SAI_DB_RESULT_BUSY; } @@ -227,7 +228,7 @@ sais_event_reset(struct vhd *vhd, const char *event_uuid) ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; @@ -236,7 +237,7 @@ sais_event_reset(struct vhd *vhd, const char *event_uuid) sqlite3_free(err); } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return SAI_DB_RESULT_OK; @@ -254,14 +255,15 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid) size_t len; int ret; - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb) == 0) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, 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); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; @@ -281,14 +283,14 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid) ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; return SAI_DB_RESULT_ERROR; } } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); } @@ -303,7 +305,7 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid) return SAI_DB_RESULT_ERROR; } - sais_event_db_delete_database(vhd, event_uuid); + sai_event_db_delete_database(vhd->sqlite3_path_lhs, event_uuid); sais_eventchange(vhd->h_ss_websrv, event_uuid, SAIES_DELETED); len = (size_t)lws_snprintf(pre + LWS_PRE, sizeof(pre) - LWS_PRE, @@ -326,14 +328,15 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid) sai_db_result_t sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) { + char filt[256], esc[96]; + struct lwsac *ac = NULL; 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)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) return SAI_DB_RESULT_ERROR; lws_sql_purify(esc, platform, sizeof(esc)); @@ -344,7 +347,7 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) &o, &ac, 0, 999) >= 0) { ret = sqlite3_exec(pdb, "BEGIN TRANSACTION", NULL, NULL, &err); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; @@ -356,7 +359,7 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) 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); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return SAI_DB_RESULT_BUSY; } @@ -364,7 +367,7 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) ret = sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); if (ret != SQLITE_OK) { - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); if (ret == SQLITE_BUSY) return SAI_DB_RESULT_BUSY; @@ -373,7 +376,7 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) sqlite3_free(err); } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return SAI_DB_RESULT_OK; diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 7f7c472..ae5b599 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -107,7 +107,8 @@ 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)) { + if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { /* * Empty the task-specific log cache into the event- @@ -125,7 +126,7 @@ sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); if (err) sqlite3_free(err); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); } else lwsl_err("%s: unable to open event-specific database\n", @@ -240,7 +241,8 @@ sais_log_to_db(struct vhd *vhd, sai_log_t *log) sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) return; lws_sql_purify(esc_uuid, log->task_uuid, sizeof(esc_uuid)); @@ -252,7 +254,7 @@ sais_log_to_db(struct vhd *vhd, sai_log_t *log) if (sai_sqlite3_statement(pdb, q, "update build_step")) lwsl_err("%s: failed to update build_step\n", __func__); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); } sai_plat_t * @@ -400,7 +402,8 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) sai_event_t *e = lws_container_of(pe, sai_event_t, list); sqlite3 *pdb = NULL; - if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { + if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) { sqlite3_stmt *sm; lws_snprintf(q, sizeof(q), @@ -422,7 +425,7 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) } sqlite3_finalize(sm); } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); } } lws_end_foreach_dll(pe); @@ -495,7 +498,8 @@ sais_process_rej(struct vhd *vhd, struct pss *pss, /* start build duration only from first step accepted */ sai_task_uuid_to_event_uuid(event_uuid, rej->task_uuid); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) break; lws_sql_purify(esc_uuid, rej->task_uuid, sizeof(esc_uuid)); @@ -530,7 +534,7 @@ sais_process_rej(struct vhd *vhd, struct pss *pss, lwsl_notice("%s: exiting, setting build_step %d\n", __func__, build_step); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); if (sais_set_task_state(vhd, rej->task_uuid, @@ -968,7 +972,8 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b * reason. */ - if (sais_event_db_ensure_open(pss->vhd, event_uuid, 0, + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pss->pdb_artifact)) { lwsl_err("%s: unable to open event-specific " "database\n", __func__); @@ -987,7 +992,7 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b NULL, lsm_schema_sq3_map_task, &o, &ac, 0, 1); if (n < 0 || !o.head) { - sais_event_db_close(vhd, &pss->pdb_artifact); + sai_event_db_close(&vhd->sqlite3_cache, &pss->pdb_artifact); lwsl_notice("%s: no task of that id\n", __func__); lwsac_free(&pss->a.ac); return -1; @@ -1273,7 +1278,7 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b lws_dll2_owner_clear(&o); lws_dll2_add_head(&metric->list, &o); - sai_dump_stderr((const char *)xbuf + LWS_PRE, used); + sai_dump_stderr(xbuf + LWS_PRE, used); if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) lwsl_warn("%s: unable to broadcast to web\n", __func__); @@ -1307,7 +1312,7 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b afail: lwsac_free(&ac); lwsac_free(&pss->a.ac); - sais_event_db_close(vhd, &pss->pdb_artifact); + sai_event_db_close(&vhd->sqlite3_cache, &pss->pdb_artifact); return -1; } @@ -1459,7 +1464,7 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsac_free(&task->ac_task_container); free(task); - sai_dump_stderr((const char *)start, w); + sai_dump_stderr(start, w); lwsl_err("%s: ########## ATTACH TASK --^\n", __func__); diff --git a/src/server/s-ws-web.c b/src/server/s-ws-web.c index cdc2d90..0ccc779 100644 --- a/src/server/s-ws-web.c +++ b/src/server/s-ws-web.c @@ -67,6 +67,13 @@ static const lws_struct_map_t lsm_viewercount_members[] = { LSM_UNSIGNED(sai_viewer_state_t, viewers, "count"), }; +static lws_struct_map_t lsm_browser_taskinfo[] = { + LSM_CARRAY (sai_browse_rx_taskinfo_t, task_hash, "task_hash"), + LSM_UNSIGNED (sai_browse_rx_taskinfo_t, logs, "logs"), + LSM_UNSIGNED (sai_browse_rx_taskinfo_t, js_api_version, "js_api_version"), + LSM_UNSIGNED (sai_browse_rx_taskinfo_t, last_log_ts, "last_log_ts"), +}; + static const lws_struct_map_t lsm_schema_json_map[] = { LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_browser_taskreset, /* shares struct */ "com.warmcat.sai.taskreset"), @@ -86,6 +93,8 @@ static const lws_struct_map_t lsm_schema_json_map[] = { "com.warmcat.sai.platreset"), LSM_SCHEMA (sai_stay_t, NULL, lsm_stay, "com.warmcat.sai.stay"), + LSM_SCHEMA (sai_browse_rx_taskinfo_t, NULL, lsm_browser_taskinfo, + "com.warmcat.sai.taskinfo") }; enum { @@ -98,6 +107,7 @@ enum { SAIS_WS_WEBSRV_RX_REBUILD, SAIS_WS_WEBSRV_RX_PLATRESET, SAIS_WS_WEBSRV_RX_STAY, + SAIS_WS_WEBSRV_RX_TASKINFO, }; static int @@ -154,10 +164,10 @@ sais_list_builders(struct vhd *vhd) subsequent = 0; unsigned int ss_flags = LWSSS_FLAG_SOM; lws_dll2_owner_t db_builders_owner; - sai_plat_t *builder_from_db; lws_struct_serialize_t *js; struct lwsac *ac = NULL; lws_wsmsg_info_t info; + sai_plat_t *sp; size_t w; memset(&db_builders_owner, 0, sizeof(db_builders_owner)); @@ -178,40 +188,48 @@ sais_list_builders(struct vhd *vhd) lws_struct_json_serialize_result_t r; sai_plat_t *live_builder; - builder_from_db = lws_container_of(walk, sai_plat_t, sai_plat_list); - live_builder = sais_builder_from_uuid(vhd, builder_from_db->name); + sp = lws_container_of(walk, sai_plat_t, sai_plat_list); + live_builder = sais_builder_from_uuid(vhd, sp->name); if (live_builder) { - builder_from_db->online = 1; - lws_strncpy(builder_from_db->peer_ip, live_builder->peer_ip, - sizeof(builder_from_db->peer_ip)); - builder_from_db->stay_on = live_builder->stay_on; + sp->online = 1; + lws_strncpy(sp->peer_ip, live_builder->peer_ip, + sizeof(sp->peer_ip)); + sp->stay_on = live_builder->stay_on; } else - builder_from_db->online = 0; + sp->online = 0; - builder_from_db->powering_up = 0; - builder_from_db->powering_down = 0; + sp->powering_up = 0; + sp->powering_down = 0; lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.power_state_owner.head) { sai_power_state_t *ps = lws_container_of(p, sai_power_state_t, list); - size_t host_len = strlen(ps->host), pl = strlen(builder_from_db->name); - - lwsl_notice("%s: %s vs %s\n", __func__, builder_from_db->name, ps->host); + size_t host_len = strlen(ps->host), pl = strlen(sp->name); - if ((!strncmp(builder_from_db->name, ps->host, host_len) && - builder_from_db->name[host_len] == '.') || (pl > host_len && - !strncmp(builder_from_db->name + (pl - host_len), ps->host, host_len))) + if ((!strncmp(sp->name, ps->host, host_len) && + sp->name[host_len] == '.') || (pl > host_len && + !strncmp(sp->name + (pl - host_len), ps->host, host_len))) { + lwsl_notice("%s: %s vs %s, sp->online %d, pup %d, pdwn %d\n", __func__, sp->name, ps->host, sp->online, ps->powering_up, ps->powering_down); + /* + * powering_up/down comes to us as a one-shot + * notification, we have to clear our copy of it + */ + if (sp->online) + ps->powering_up = 0; + else + ps->powering_down = 0; + lwsl_notice("%s: adjusting powering_ %d %d\n", __func__, ps->powering_up, ps->powering_down); - builder_from_db->powering_up = ps->powering_up; - builder_from_db->powering_down = ps->powering_down; + sp->powering_up = ps->powering_up; + sp->powering_down = ps->powering_down; break; } } lws_end_foreach_dll(p); js = lws_struct_json_serialize_create(lsm_schema_map_plat_simple, LWS_ARRAY_SIZE(lsm_schema_map_plat_simple), - 0, builder_from_db); + 0, sp); if (!js) goto bail; @@ -300,6 +318,7 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) memset(&a, 0, sizeof(a)); a.map_st[0] = lsm_schema_json_map; + a.map_st[1] = lsm_schema_json_map; a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_json_map); a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_json_map); a.ac_block_size = 128; @@ -495,71 +514,9 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, int *flags) { websrvss_srv_t *m = (websrvss_srv_t *)userobj; - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->bl_srv_to_web), depi, fl; - char som, som1, eom, final = 1; - size_t fsl, used; - - if (!m->bl_srv_to_web) - return LWSSSSRET_TX_DONT_SEND; - - depi = *pi; - - /* - * We can only issue *len at a time. - * - * Notice we are getting the stored flags from the START of the fragment each time. - * that means we can still see the right flags stored with the fragment, even if we - * have partially used the buflist frag and are partway through it. - * - * Ergo, only something to skip if we are at som=1. And also notice that although - * *pi will be right, after the lws_buflist..._use() api, what it points to has been - * destroyed. So we also dereference *pi into depi for use below. - */ - - fsl = lws_buflist_next_segment_len(&m->bl_srv_to_web, NULL); - - lws_buflist_fragment_use(&m->bl_srv_to_web, NULL, 0, &som, &eom); - if (som) { - fsl -= sizeof(int); - lws_buflist_fragment_use(&m->bl_srv_to_web, buf, sizeof(int), &som1, &eom); - } - if (!(depi & LWSSS_FLAG_SOM)) - som = 0; - - used = (size_t)lws_buflist_fragment_use(&m->bl_srv_to_web, (uint8_t *)buf, *len, &som1, &eom); - if (!used) - return LWSSSSRET_TX_DONT_SEND; - - if (used < fsl || !(depi & LWSSS_FLAG_EOM)) /* we saved SS flags at the start of the buf */ - final = 0; - - *len = used; - fl = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - - // lwsl_ss_notice(m->ss, "Sending %d srv->web: ssflags %d", (int)*len, fl); - if ((fl & LWSSS_FLAG_SOM) && (((*flags) & 3) == 2)) { - lwsl_ss_err(m->ss, "TX: Illegal LWSSS_FLAG_SOM after previous frame without LWSSS_FLAG_EOM"); - assert(0); - } - if (!(fl & LWSSS_FLAG_SOM) && ((*flags) & 3) == 3) { - lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM after previous frame with LWSSS_FLAG_EOM"); - assert(0); - } - if (!(fl & LWSSS_FLAG_SOM) && !((*flags) & 2)) { - lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM on first frame"); - assert(0); - } - - *flags = fl; - - - // lwsl_hexdump_notice(buf, *len); - - if (m->bl_srv_to_web) - return lws_ss_request_tx(m->ss); - - return 0; + return sai_ss_tx_from_buflist_helper(m->ss, &m->bl_srv_to_web, + buf, len, flags); } diff --git a/src/web/CMakeLists.txt b/src/web/CMakeLists.txt index 481ebf7..cd167a5 100644 --- a/src/web/CMakeLists.txt +++ b/src/web/CMakeLists.txt @@ -24,6 +24,7 @@ if (requirements) if (APPLE) set_property(TARGET sai-web PROPERTY MACOSX_RPATH YES) endif() + # set(CMAKE_C_FLAGS "-ffunction-sections -Wl,--gc-sections -Wl,--print-gc-sections ${CMAKE_C_FLAGS}" ) # # sqlite3 paths (web) diff --git a/src/web/w-artifact.c b/src/web/w-artifact.c index 553a2a9..d5007fb 100644 --- a/src/web/w-artifact.c +++ b/src/web/w-artifact.c @@ -83,7 +83,8 @@ saiw_get_blob(struct vhd *vhd, const char *url, sqlite3 **pdb, /* open the event-specific database object */ - if (sais_event_db_ensure_open(vhd, event_uuid, 0, pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, pdb)) { lwsl_info("%s: unable to open event-specific database\n", __func__); @@ -143,7 +144,7 @@ saiw_get_blob(struct vhd *vhd, const char *url, sqlite3 **pdb, fail: lwsac_free(&ac); - sais_event_db_close(vhd, pdb); + sai_event_db_close(&vhd->sqlite3_cache, pdb); lwsl_notice("%s: couldn't find blob %s\n", __func__, url); diff --git a/src/web/w-comms.c b/src/web/w-comms.c index aee5e74..d113ef8 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -62,202 +62,6 @@ const lws_struct_map_t lsm_schema_sq3_map_auth[] = { extern const lws_struct_map_t lsm_schema_sq3_map_event[]; -/* 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-web) 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 1; - } - - /* 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 for %s\n", __func__, filepath); - return 1; - } - - 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 for %s\n", __func__, filepath); - - return 1; - } - - sai_sqlite3_statement(*ppdb, - "CREATE INDEX IF NOT EXISTS logs_index ON logs (task_uuid, timestamp);", - "create logs index"); - - if (lws_struct_sq3_create_table(*ppdb, lsm_schema_sq3_map_artifact)) { - lwsl_err("%s: unable to create artifact table for %s\n", __func__, filepath); - - return 1; - } - - 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 1; - } - - 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, m-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_delete_database(struct vhd *vhd, const char *event_uuid) -{ - char filepath[256], saf[33], 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); - - if (unlink(filepath)) { - lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, - filepath, errno); - ra = 1; - } - - lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3-wal", - vhd->sqlite3_path_lhs, saf); - - if (unlink(filepath)) { - lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, - filepath, errno); - ra = 1; - } - - lws_snprintf(filepath, sizeof(filepath), "%s-event-%s.sqlite3-shm", - vhd->sqlite3_path_lhs, saf); - - if (unlink(filepath)) { - lwsl_err("%s (web): unable to delete %s (%d)\n", __func__, - filepath, errno); - ra = 1; - } - - if (!ra) - lwsl_notice("%s (web): deleted %s OK\n", __func__, filepath); - - return ra; -} - typedef enum { SHMUT_NONE = -1, diff --git a/src/web/w-private.h b/src/web/w-private.h index 64e18fa..eff812c 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -167,14 +167,6 @@ struct pss { unsigned int toggle_favour_sch:1; }; -typedef struct sais_sqlite_cache { - lws_dll2_t list; - char uuid[65]; - sqlite3 *pdb; - lws_usec_t idle_since; - int refcount; -} sais_sqlite_cache_t; - struct vhd { struct lws_context *context; struct lws_vhost *vhost; @@ -205,6 +197,16 @@ struct vhd { lws_sorted_usec_list_t sul_logcache; }; +typedef struct saiw_websrv { + struct lws_ss_handle *ss; + void *opaque_data; + + lws_struct_args_t a; + struct lejp_ctx ctx; + struct lws_buflist *wbltx; +} saiw_websrv_t; + + extern struct lws_context * sai_lws_context_from_json(const char *config_dir, struct lws_context_creation_info *info, @@ -219,18 +221,6 @@ sai_notification_file_upload_cb(void *data, const char *name, enum lws_spa_fileupload_states state); int -sai_sqlite3_statement(sqlite3 *pdb, const char *cmd, const char *desc); - -int -sais_event_db_ensure_open(struct vhd *vhd, const char *event_uuid, char can_create, sqlite3 **ppdb); - -void -sais_event_db_close(struct vhd *vhd, sqlite3 **ppdb); - -int -sais_event_db_delete_database(struct vhd *vhd, const char *event_uuid); - -int sai_sq3_event_lookup(sqlite3 *pdb, uint64_t start, lws_struct_args_cb cb, void *ca); int @@ -268,9 +258,6 @@ int saiw_task_cancel(struct vhd *vhd, const char *task_uuid); int -saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags); - -int saiw_get_blob(struct vhd *vhd, const char *url, sqlite3 **pdb, sqlite3_blob **blob, uint64_t *length); @@ -288,7 +275,8 @@ saiw_sched_destroy(struct lws_dll2 *d, void *user); void -saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int min_api_version, enum lws_write_protocol flags); +saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, + enum lws_write_protocol flags); void saiw_browser_state_changed(struct pss *pss, int established); diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 07b1397..0d0d625 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -289,7 +289,8 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) /* open the event-specific database object */ - if (sais_event_db_ensure_open(pss->vhd, event_uuid, 0, &pdb)) { + 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); return 0; @@ -305,7 +306,7 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) 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); - sais_event_db_close(pss->vhd, &pdb); + sai_event_db_close(&pss->vhd->sqlite3_cache, &pdb); if (n < 0 || !o.head) goto bail; @@ -492,7 +493,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, if (saiw_pss_schedule_taskinfo(pss, ti->task_hash, !!ti->logs)) goto soft_error; - break; + goto ok; case SAIM_WS_BROWSER_RX_EVENTINFO: @@ -501,7 +502,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, if (saiw_pss_schedule_eventinfo(pss, ei->event_hash)) goto soft_error; - break; + goto ok; case SAIM_WS_BROWSER_RX_TASKRESET: @@ -513,8 +514,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, */ ei = (sai_browse_rx_evinfo_t *)a.dest; - - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_STAY: @@ -528,8 +527,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, /* * User is asking us to set or release a stay on a builder */ - - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_TASKREBUILDLASTSTEP: @@ -541,8 +538,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, */ ei = (sai_browse_rx_evinfo_t *)a.dest; - - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_EVENTRESET: @@ -558,8 +553,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsl_notice("%s: received request to reset event %s\n", __func__, ei->event_hash); - - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_EVENTDELETE: @@ -575,9 +568,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, lwsl_notice("%s: received request to delete event %s\n", __func__, ei->event_hash); - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); - lwsac_free(&a.ac); - break; case SAIM_WS_BROWSER_RX_TASKCANCEL: @@ -595,7 +585,7 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, __func__, can->task_uuid); saiw_task_cancel(vhd, can->task_uuid); - break; + goto ok; case SAIM_WS_BROWSER_RX_REBUILD: if (!sais_conn_auth(pss)) @@ -604,8 +594,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, /* * User is asking us to rebuild a builder */ - - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; case SAIM_WS_BROWSER_RX_PLATRESET: @@ -615,8 +603,6 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, /* * User is asking us to reset / rebuild a whole platform */ - - saiw_websrv_queue_tx(vhd->h_ss_websrv, buf, bl, ss_flags); break; default: @@ -624,6 +610,11 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, break; } + sai_ss_queue_frag_on_buflist_REQUIRES_LWS_PRE(vhd->h_ss_websrv, + &((saiw_websrv_t *)lws_ss_to_user_object(vhd->h_ss_websrv))->wbltx, + buf, bl, ss_flags); + +ok: ret = 0; bail: @@ -686,8 +677,6 @@ again: // lwsl_notice("%s: send_state %d, pss %p, wsi %p\n", __func__, // pss->send_state, pss, pss->wsi); - // 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); @@ -729,7 +718,8 @@ again: lwsl_info("%s: collecting logs %s\n", __func__, esc); - if (sais_event_db_ensure_open(vhd, event_uuid, 0, + 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__); @@ -743,7 +733,7 @@ again: &pss->logs_owner, &pss->logs_ac, 0, 100); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); if (sr) { @@ -1033,7 +1023,8 @@ enum_tasks: do { task_ac = NULL; - if (sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { + if (sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) { lwsl_err("%s: unable to open event-specific database\n", __func__); @@ -1045,11 +1036,11 @@ enum_tasks: lsm_schema_sq3_map_task, &task_owner, &task_ac, sch->task_index, 1)) { lwsl_err("%s: OVERVIEW 1 failed\n", __func__); - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); break; } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); if (!task_owner.count) break; @@ -1265,7 +1256,8 @@ b_finish: // __func__, event_uuid); lws_dll2_owner_clear(&sch->owner); - if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + 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); @@ -1281,7 +1273,7 @@ b_finish: &sch->ac, 0, 10)) { lwsl_err("%s: get afcts failed\n", __func__); } - sais_event_db_close(vhd, &pdb); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); } } diff --git a/src/web/w-ws-server.c b/src/web/w-ws-server.c index 20527de..8bf2bae 100644 --- a/src/web/w-ws-server.c +++ b/src/web/w-ws-server.c @@ -32,15 +32,6 @@ #include "w-private.h" -typedef struct saiw_websrv { - struct lws_ss_handle *ss; - void *opaque_data; - - lws_struct_args_t a; - struct lejp_ctx ctx; - struct lws_buflist *wbltx; -} saiw_websrv_t; - static lws_struct_map_t lsm_websrv_evinfo[] = { LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"), }; @@ -81,15 +72,13 @@ enum { * The flags are lws_write() flags. */ void -saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int api_ver_min, enum lws_write_protocol flags) +saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, + enum lws_write_protocol flags) { - int eff = 0; - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { struct pss *pss = lws_container_of(p, struct pss, same); int *pi = (int *)((const char *)buf - sizeof(int)); - eff++; *pi = (int)flags; if (lws_buflist_append_segment(&pss->raw_tx, buf - sizeof(int), len + sizeof(int)) < 0) @@ -100,31 +89,6 @@ saiw_ws_broadcast_raw(struct vhd *vhd, const void *buf, size_t len, unsigned int } lws_end_foreach_dll(p); } - -/* - * Queue messages to send from sai-web to sai-server - */ - -int -saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) -{ - saiw_websrv_t *m = (saiw_websrv_t *)lws_ss_to_user_object(h); - unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); - - *pi = ss_flags; - - // lwsl_ss_notice(h, "sai-web: Queuing sai-web -> sai-server"); - // lwsl_hexdump_notice(buf, len); - - if (lws_buflist_append_segment(&m->wbltx, buf - sizeof(int), len + sizeof(int)) < 0) - lwsl_ss_err(h, "failed to append"); /* still ask to drain */ - - if (lws_ss_request_tx(h)) - lwsl_ss_err(h, "failed to request tx"); - - return 0; -} - /* * sai-web is receiving from sai-server * @@ -175,17 +139,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, 0, + saiw_ws_broadcast_raw(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, 0, + saiw_ws_broadcast_raw(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, 0, + saiw_ws_broadcast_raw(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); @@ -201,7 +165,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, 0, + saiw_ws_broadcast_raw(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); @@ -277,11 +241,11 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) case SAIS_WS_WEBSRV_RX_LOADREPORT: // lwsl_notice("%s: ^^^^^^^^^^^^^^ SAIS_WS_WEBSRV_RX_LOADREPORT forwarding to browser\n", __func__); // lwsl_hexdump_notice(buf, len); - saiw_ws_broadcast_raw(vhd, buf, len - (unsigned int)n, 0, + saiw_ws_broadcast_raw(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, 0, + saiw_ws_broadcast_raw(vhd, buf, len - (unsigned int)n, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); break; } @@ -306,55 +270,8 @@ saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, int *flags) { saiw_websrv_t *m = (saiw_websrv_t *)userobj; - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->wbltx), depi; - char som, som1, eom, final = 1; - size_t fsl, used; - if (!m->wbltx) { - // lwsl_notice("%s: nothing to send from web -> srv\n", __func__); - return LWSSSSRET_TX_DONT_SEND; - } - - depi = *pi; - - /* - * We can only issue *len at a time. - * - * Notice we are getting the stored flags from the START of the fragment each time. - * that means we can still see the right flags stored with the fragment, even if we - * have partially used the buflist frag and are partway through it. - * - * Ergo, only something to skip if we are at som=1. And also notice that although - * *pi will be right, after the lws_buflist..._use() api, what it points to has been - * destroyed. So we also dereference *pi into depi for use below. - */ - - fsl = lws_buflist_next_segment_len(&m->wbltx, NULL); - - lws_buflist_fragment_use(&m->wbltx, NULL, 0, &som, &eom); - if (som) { - fsl -= sizeof(int); - lws_buflist_fragment_use(&m->wbltx, buf, sizeof(int), &som1, &eom); - } - - /* this is the only buflist user on pss->raw_tx */ - used = (size_t)lws_buflist_fragment_use(&m->wbltx, (uint8_t *)buf, *len, &som1, &eom); - if (!used) - return LWSSSSRET_TX_DONT_SEND; - - if (used < fsl || (depi & LWS_WRITE_NO_FIN)) - final = 0; - - *len = used; - *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - - // lwsl_ss_notice(m->ss, "Sending %d web->srv: ssflags %d", (int)*len, (int)*flags); - // lwsl_hexdump_notice(buf, *len); - - if (m->wbltx) - return lws_ss_request_tx(m->ss); - - return 0; + return sai_ss_tx_from_buflist_helper(m->ss, &m->wbltx, buf, len, flags); } static int @@ -441,5 +358,7 @@ saiw_update_viewer_count(struct vhd *vhd) 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); + sai_ss_queue_frag_on_buflist_REQUIRES_LWS_PRE(vhd->h_ss_websrv, + &((saiw_websrv_t *)lws_ss_to_user_object(vhd->h_ss_websrv))->wbltx, + buf + LWS_PRE, len, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); }
Page fetched 0s ago, creation time: 19ms (vhost etag hits: 0%, cache hits: 0%)