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