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);
}