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