diff --git a/assets/arch-x86_64.svg b/assets/arch-x86_64.svg
new file mode 100644
index 0000000..39eaa23
--- /dev/null
+++ b/assets/arch-x86_64.svg
@@ -0,0 +1 @@
+<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 3.816 3.816" height="14.423" width="14.423" xmlns:v="https://vecta.io/nano"><circle cx="1.908" cy="1.908" r="1.908" opacity=".12" fill="#f9f9f9"/><path d="M1.285 1.4l-.607.59v1.15H1.92l.528-.572-1.155.018zM.735.784h2.327v2.36l-.594-.594V1.382H1.296z" paint-order="markers fill stroke" fill="#099b39"/></svg>
diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c
index b6b3ab0..e0d7992 100644
--- a/src/builder/b-nspawn.c
+++ b/src/builder/b-nspawn.c
@@ -74,6 +74,7 @@ saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel)
((int)ns->sp->nspawn_owner.count - 1),
saib_get_free_ram_kib(),
saib_get_free_disk_kib(builder.home));
+ ns->retcode_set = 0;
}
n += lws_snprintf(lj + n, sizeof(lj) - (unsigned int)n,
@@ -110,7 +111,7 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason,
lwsl_info("%s: stdwsi CLOSE, ns %p, lsp: %p, wsi: %p, fd: %d, stdfd: %d\n",
__func__, op ? op->ns : NULL, op ? op->lsp : NULL,
wsi, lws_get_socket_fd(wsi), lws_spawn_get_stdfd(wsi));
-
+/*
ilen = lws_snprintf((char *)buf, sizeof(buf), "Stdwsi %d close\n", lws_spawn_get_stdfd(wsi));
if (ns) {
saib_log_chunk_create(ns, buf, (size_t)ilen, 3);
@@ -119,7 +120,7 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason,
lwsl_warn("%s: lws_ss_request_tx failed\n",
__func__);
}
-
+*/
if (op && op->lsp) {
if (lws_spawn_stdwsi_closed(op->lsp, wsi) &&
ns->reap_cb_called) {
@@ -400,7 +401,7 @@ static const char * const runscript_win_first =
"set SAI_LOGPROXY_TTY0=%s\n"
"set SAI_LOGPROXY_TTY1=%s\n"
"set HOME=%s\n"
- "cd %s%s &&"
+ "cd %s%s &&"
" rmdir /s /q src & "
"%s < NUL"
;
@@ -413,14 +414,14 @@ static const char * const runscript_win_next =
"set SAI_LOGPROXY_TTY0=%s\n"
"set SAI_LOGPROXY_TTY1=%s\n"
"set HOME=%s\n"
- "cd %s%s &&"
+ "cd %s%s &&"
"%s < NUL"
;
#else
static const char * const runscript_first =
- "#!/bin/bash -x\n"
+ "#!/bin/bash\n" /* use -x to see what it does for these */
#if defined(__APPLE__)
"export PATH=/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/sbin:/usr/sbin\n"
#else
@@ -445,7 +446,7 @@ static const char * const runscript_first =
;
static const char * const runscript_next =
- "#!/bin/bash -x\n"
+ "#!/bin/bash\n" /* use -x to see what it does for these */
#if defined(__APPLE__)
"export PATH=/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/sbin:/usr/sbin\n"
#else
@@ -469,7 +470,7 @@ static const char * const runscript_next =
;
static const char * const runscript_build =
- "#!/bin/bash -x\n"
+ "#!/bin/bash\n" /* use -x to see what it does for these */
#if defined(__APPLE__)
"export PATH=/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/sbin:/usr/sbin\n"
#else
diff --git a/src/builder/b-private.h b/src/builder/b-private.h
index 3d3ab5c..3bd114b 100644
--- a/src/builder/b-private.h
+++ b/src/builder/b-private.h
@@ -218,10 +218,6 @@ int
saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel);
int
-saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm,
- const char *rej_task_uuid);
-
-int
rm_rf_cb(const char *dirpath, void *user, struct lws_dir_entry *lde);
extern const struct lws_protocols protocol_logproxy, protocol_resproxy;
diff --git a/src/builder/b-task.c b/src/builder/b-task.c
index 7015f08..41ac815 100644
--- a/src/builder/b-task.c
+++ b/src/builder/b-task.c
@@ -279,9 +279,9 @@ saib_set_ns_state(struct sai_nspawn *ns, int state)
* update all servers we're connected to about builder status / optional reject
*/
-int
+static int
saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm,
- const char *rej_task_uuid)
+ const char *rej_task_uuid, unsigned int reason)
{
struct sai_rejection rej;
@@ -294,29 +294,23 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm,
memset(&rej, 0, sizeof(rej));
/*
- * Queue a builder status update /
- * optional task rejection
+ * Queue a builder task status update
*/
- if (rej_task_uuid) {
- lwsl_notice("%s: builder %s occupied reject\n",
- __func__, sp->name);
-
+ if (rej_task_uuid)
lws_strncpy(rej.task_uuid, rej_task_uuid,
sizeof(rej.task_uuid));
- } else
- lwsl_notice("%s: issuing load update\n", __func__);
- lws_snprintf(rej.host_platform, sizeof(rej.host_platform), "%s",
- sp->name);
+ lws_snprintf(rej.host_platform, sizeof(rej.host_platform), "%s", sp->name);
- rej.avail_slots = (int)(sp->job_limit ? sp->job_limit : 6u) -
+ rej.avail_slots = (int)(sp->job_limit ? sp->job_limit : 6u) -
(int)sp->nspawn_owner.count;
- rej.avail_mem_kib = saib_get_free_ram_kib();
- rej.avail_sto_kib = saib_get_free_disk_kib(builder.home);
+ rej.avail_mem_kib = saib_get_free_ram_kib();
+ rej.avail_sto_kib = saib_get_free_disk_kib(builder.home);
+
+ rej.reason = (uint8_t)reason;
- if (saib_srv_queue_json_fragments_helper(spm->ss,
- lsm_schema_json_task_rej,
+ if (saib_srv_queue_json_fragments_helper(spm->ss, lsm_schema_json_task_rej,
LWS_ARRAY_SIZE(lsm_schema_json_task_rej), &rej))
return -1;
@@ -381,7 +375,9 @@ saib_task_destroy(struct sai_nspawn *ns)
* Schedule informing all the servers we're connected to
*/
- saib_queue_task_status_update(ns->sp, ns->spm, NULL);
+ if (ns->task)
+ saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid,
+ SAI_TASK_REASON_DESTROYED);
}
if (ns->task && ns->task->ac_task_container) {
@@ -736,12 +732,13 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len)
}
switch (a.top_schema_index) {
+
case SAIB_RX_TASK_ALLOCATION:
task = (sai_task_t *)a.dest;
task->ac_task_container = a.ac; /* bequeath lwsac responsibility */
/*
- * Master is requesting that a platform adopt a task...
+ * Server is requesting that a platform adopt a task...
*
* Multiple platforms may be using this connection to a given
* server so we have to disambiguate which platform he's
@@ -764,20 +761,11 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len)
return 1;
}
- /*
- * There's not already an existing step we accepted,
- * using the same uuid?
- */
-
lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, sp->nspawn_owner.head) {
struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list);
lwsl_notice("%s: nspawn_census: %s\n", __func__, xns->task->uuid);
-// if (xns->task && !strcmp(xns->task->uuid, task->uuid)) {
-// lwsl_err("%s: server offered task %s that already has an extant nspawn\n", __func__, task->uuid);
- // return 0;
-// }
} lws_end_foreach_dll_safe(d, d1);
@@ -790,7 +778,7 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len)
sp->deserialization_ac = a.ac;
/*
- * Look for a spare nspawn...
+ * Are we willing to take this task step on?
*
* We may connect to multiple servers and it's asynchronous
* which server may have tasked us first, so it's not that
@@ -804,15 +792,17 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len)
struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list);
if (xns->task && !strcmp(xns->task->uuid, task->uuid)) {
lwsl_warn("%s: server offered task that's already running\n", __func__);
- saib_queue_task_status_update(sp, spm, task->uuid);
+ saib_queue_task_status_update(sp, spm, task->uuid, SAI_TASK_REASON_DUPE);
return 0;
}
if (!xns->task && !ns)
ns = xns;
} lws_end_foreach_dll_safe(d, d1);
- if (saib_can_accept_task(task, sp)) { /* not accepted */
- if (saib_queue_task_status_update(sp, spm, task->uuid))
+
+ if (saib_can_accept_task(task, sp)) {
+ lwsl_warn("%s: builder rejects offered task\n", __func__);
+ if (saib_queue_task_status_update(sp, spm, task->uuid, SAI_TASK_REASON_BUSY))
return -1;
return 0;
}
@@ -1061,17 +1051,15 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len)
}
/* we're busy, we're not in the mood for suspending */
+
lwsl_notice("%s: cancelling suspend grace time\n", __func__);
lws_sul_cancel(&ns->builder->sul_idle);
/*
- * Let the mirror thread get on with things...
- *
- * When we took on a task, we should inform any servers we're
- * connected to about our change in task load status
+ * We accepted the task
*/
- if (saib_queue_task_status_update(sp, spm, NULL))
+ if (saib_queue_task_status_update(sp, spm, task->uuid, SAI_TASK_REASON_ACCEPTED))
goto bail;
break;
diff --git a/src/common/include/private.h b/src/common/include/private.h
index aead6b1..a55abd8 100644
--- a/src/common/include/private.h
+++ b/src/common/include/private.h
@@ -106,10 +106,14 @@ typedef struct sai_platform_load {
struct sai_nspawn;
+typedef struct sai_plat sai_plat_t;
+
typedef struct {
- lws_dll2_t list;
+ lws_dll2_t list; /* managed by an owner via LSM_SCHEMA_DLL2 / lsm_task */
lws_dll2_t pending_assign_list;
+
const struct sai_event *one_event; /* event we are associated with */
+
char platform[96];
char build[4096]; /* strsubst and serialized */
char taskname[96];
@@ -129,11 +133,11 @@ typedef struct {
struct lwsac *ac_task_container;
- const char *server_name; /* used in offer */
- const char *repo_name; /* used in offer */
- const char *git_ref; /* used in offer */
- const char *git_hash; /* used in offer */
- const char *git_repo_url; /* used in offer */
+ const char *server_name; /* used in offer */
+ const char *repo_name; /* used in offer */
+ const char *git_ref; /* used in offer */
+ const char *git_hash; /* used in offer */
+ const char *git_repo_url; /* used in offer */
uint64_t last_updated;
uint64_t started;
uint64_t duration;
@@ -153,8 +157,6 @@ typedef struct {
char rebuildable;
} sai_task_t;
-typedef struct sai_plat sai_plat_t;
-
struct saib_logproxy {
char sockpath[128];
struct sai_nspawn *ns;
@@ -228,15 +230,25 @@ struct sai_nspawn {
/*
* Builder is indicating he can't take the task and server should free it up
* and try another builder.
+ *
*/
+enum {
+ SAI_TASK_REASON_ACCEPTED = 0,
+ SAI_TASK_REASON_DUPE = 1,
+ SAI_TASK_REASON_BUSY = 2,
+ SAI_TASK_REASON_DESTROYED = 3,
+};
+
typedef struct sai_rejection {
struct lws_dll2 list;
+
char host_platform[65];
char task_uuid[65];
int avail_slots;
unsigned int avail_mem_kib;
unsigned int avail_sto_kib;
+ unsigned char reason;
} sai_rejection_t;
/*
@@ -268,7 +280,6 @@ struct sai_event;
typedef struct sai_event {
struct lws_dll2 list;
- lws_dll2_owner_t task_owner;
char repo_name[65];
char repo_fetchurl[96];
char ref[65];
@@ -425,8 +436,9 @@ 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;
- char uuid[65];
+ lws_dll2_t list;
+ char uuid[65];
+ char started;
} sai_uuid_list_t;
/*
@@ -608,7 +620,7 @@ extern const lws_struct_map_t
lsm_artifact[8],
lsm_plat_list[1],
lsm_schema_map_plat[1],
- lsm_task_rej[5],
+ lsm_task_rej[6],
lsm_task_cancel[1],
lsm_schema_json_map_can[1],
lsm_schema_json_map_task[1],
diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c
index 52b3bb0..9fb2d92 100644
--- a/src/common/struct-metadata.c
+++ b/src/common/struct-metadata.c
@@ -214,6 +214,7 @@ const lws_struct_map_t lsm_task_rej[] = {
LSM_JO_SIGNED (sai_rejection_t, avail_slots, "avail_slots"),
LSM_JO_UNSIGNED (sai_rejection_t, avail_mem_kib, "avail_mem_kib"),
LSM_JO_UNSIGNED (sai_rejection_t, avail_sto_kib, "avail_sto_kib"),
+ LSM_JO_UNSIGNED (sai_rejection_t, reason, "reason"),
};
const lws_struct_map_t lsm_schema_json_task_rej[] = {
diff --git a/src/server/s-central.c b/src/server/s-central.c
index 8f87585..cecbd6f 100644
--- a/src/server/s-central.c
+++ b/src/server/s-central.c
@@ -165,8 +165,7 @@ sais_central_cb(lws_sorted_usec_list_t *sul)
* try to bind outstanding task to specific builder
* instance
*/
- sais_allocate_task(vhd,
- (struct pss *)lws_wsi_user(cb->wsi),
+ sais_allocate_task(vhd, (struct pss *)lws_wsi_user(cb->wsi),
cb, cb->platform);
} lws_end_foreach_dll(p);
diff --git a/src/server/s-comms.c b/src/server/s-comms.c
index 92097a4..fc4fdb4 100644
--- a/src/server/s-comms.c
+++ b/src/server/s-comms.c
@@ -755,7 +755,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user,
case LWS_CALLBACK_CLOSED:
lwsac_free(&pss->query_ac);
- lwsl_wsi_user(wsi, "############################### sai-server: CLOSED sai-web conn ###############");
+ lwsl_wsi_user(wsi, "############################### sai-server: CLOSED builder conn ###############");
/* remove pss from vhd->builders (active connection list) */
lws_dll2_remove(&pss->same);
@@ -773,6 +773,21 @@ 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
@@ -856,15 +871,15 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user,
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);
+ // 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);
+ } // 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);
@@ -909,7 +924,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user,
/*
* 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));
+ // 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))
return -1;
diff --git a/src/server/s-private.h b/src/server/s-private.h
index 7b9e39f..d16ff85 100644
--- a/src/server/s-private.h
+++ b/src/server/s-private.h
@@ -123,10 +123,6 @@ struct pss {
lws_dll2_owner_t viewer_state_owner;
lws_struct_args_t a;
- union {
- sai_plat_t *b;
- sai_plat_owner_t *o;
- } u;
const char *server_name;
struct lwsac *query_ac;
@@ -352,3 +348,12 @@ sais_mark_all_builders_offline(struct vhd *vhd);
int
sql3_get_string_cb(void *user, int cols, char **values, char **name);
+int
+sais_is_task_inflight(struct vhd *vhd, const char *uuid, sai_uuid_list_t **hit);
+
+int
+sais_add_to_inflight_list_if_absent(struct vhd *vhd, sai_plat_t *sp, const char *uuid);
+
+void
+sais_inflight_entry_destroy(sai_uuid_list_t *ul);
+
diff --git a/src/server/s-task.c b/src/server/s-task.c
index 5b4037b..e13cec3 100644
--- a/src/server/s-task.c
+++ b/src/server/s-task.c
@@ -329,6 +329,77 @@ sais_event_ran_platform(struct vhd *vhd, const char *event_uuid,
}
/*
+ * On the server's builder-platform, we keep a list of tasks we have offered it.
+ *
+ * If the builder accepted the task, then we change the task's state in sqlite and
+ * remove it from this list.
+ *
+ * Inbetweentimes, we know to avoid re-offering or cancelling the task by seeing
+ * if the task is already listed as "inflight".
+ */
+
+int
+sais_is_task_inflight(struct vhd *vhd, const char *uuid, sai_uuid_list_t **hit)
+{
+
+ /*
+ * lookup a uuid across all builder / plats
+ * to see if it is inflight
+ */
+
+ 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(struct lws_dll2 *, pif,
+ build->inflight_owner.head) {
+ sai_uuid_list_t *ul = lws_container_of(pif, sai_uuid_list_t, list);
+
+ lwsl_notice("%s: '%s' vs '%s'\n", __func__, uuid, ul->uuid);
+ if (!strcmp(uuid, ul->uuid)) {
+ if (hit)
+ *hit = ul;
+ return 1;
+ }
+
+ } lws_end_foreach_dll(pif);
+ } lws_end_foreach_dll(pb);
+
+ return 0;
+}
+
+int
+sais_add_to_inflight_list_if_absent(struct vhd *vhd, sai_plat_t *sp, const char *uuid)
+{
+ sai_uuid_list_t *uuid_list;
+
+ if (sais_is_task_inflight(vhd, uuid, NULL))
+ return 0;
+
+ uuid_list = malloc(sizeof(*uuid_list));
+ if (!uuid_list)
+ return 1;
+
+ memset(uuid_list, 0, sizeof(*uuid_list));
+ lws_strncpy(uuid_list->uuid, uuid, sizeof(uuid_list->uuid));
+
+ lws_dll2_add_tail(&uuid_list->list, &sp->inflight_owner);
+
+ lwsl_notice("%s: ### created uuid_list entry for %s\n", __func__, uuid_list->uuid);
+ assert(sais_is_task_inflight(vhd, uuid, NULL));
+ return 0;
+}
+
+void
+sais_inflight_entry_destroy(sai_uuid_list_t *ul)
+{
+ lwsl_notice("%s: ### REMOVING uuid_list entry for %s\n", __func__, ul->uuid);
+
+ lws_dll2_remove(&ul->list);
+ free(ul);
+}
+
+/*
* Find the most recent task that still needs doing for platform, on any event
*/
static const sai_task_t *
@@ -383,6 +454,11 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) {
+ /*
+ * Find out how many tasks in startable state for this platform,
+ * on this event
+ */
+
lws_snprintf(query, sizeof(query), "select count(state) from tasks where "
"state = 0 and platform = '%s'", esc_plat);
m = sqlite3_exec(pdb, query, sql3_get_integer_cb, &pending_count, NULL);
@@ -392,7 +468,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
lwsl_err("%s: query failed: %d\n", __func__, m);
}
- if (pending_count > 0) {
+ if (pending_count > 0) { /* there are some startable tasks on this event */
lws_sql_purify(esc_repo, e->repo_name, sizeof(esc_repo));
lws_sql_purify(esc_ref, e->ref, sizeof(esc_ref));
@@ -410,16 +486,24 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
if (sqlite3_prepare_v2(vhd->server.pdb, query, -1, &sm, NULL) != SQLITE_OK)
break;
if (sqlite3_step(sm) == SQLITE_ROW) {
- const unsigned char *u = sqlite3_column_text(sm, 0);
- if (u)
+ const char *u = (const char *)sqlite3_column_text(sm, 0);
+ if (u) {
+ if (sais_is_task_inflight(vhd, u, NULL)) { /* we have it in hand */
+ lwsl_notice("%s: skipping pending task %s due to being inflight\n", __func__, u);
+ sqlite3_finalize(sm);
+ break;
+ }
lws_strncpy(prev_event_uuid, (const char *)u, sizeof(prev_event_uuid));
+ }
last_created = (uint64_t)sqlite3_column_int64(sm, 1);
}
sqlite3_finalize(sm);
+
if (!prev_event_uuid[0])
break;
if (!sais_event_ran_platform(vhd, prev_event_uuid, esc_plat))
continue;
+
lws_strncpy(checked_uuid, prev_event_uuid, sizeof(checked_uuid));
break;
} while (1);
@@ -504,7 +588,7 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
lsm_schema_sq3_map_task,
&owner, &pss->ac_alloc_task, 0, 1);
if (owner.count) {
- lwsl_notice("%s: MATCH! Prioritizing failed task for %s ('%s')\n",
+ lwsl_notice("%s: Prioritizing failed task for %s ('%s')\n",
__func__, platform, fti->taskname);
sais_event_db_close(vhd, &pdb);
lwsac_free(&ac);
@@ -1017,19 +1101,19 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
const sai_task_t *task_template;
char original_rejected_uuid[65];
sai_task_t temp_task;
- sai_uuid_list_t *sul;
int attempts = 0;
+#if 0
if (cb->avail_slots <= 0) {
- lwsl_info("%s: builder %s has no available slots\n", __func__,
+ lwsl_warn("%s: builder %s has no available slots\n", __func__,
cb->name);
return 1;
}
-
+#endif
lws_strncpy(original_rejected_uuid, cb->last_rej_task_uuid,
sizeof(original_rejected_uuid));
- while (attempts++ < 10) {
+ while (attempts++ < 4) {
/*
* Look for a task for this platform, on any event that needs building
@@ -1063,6 +1147,11 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
continue;
}
+ if (sais_is_task_inflight(vhd, task_template->uuid, NULL)) {
+ lwsl_notice("%s: skipping %s as listed on inflight\n", __func__, task_template->uuid);
+ continue;
+ }
+
lwsl_notice("%s: %s: task %s found for %s\n", __func__,
platform_name, task_template->uuid, cb->name);
@@ -1072,14 +1161,10 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
SAIES_PASSED_TO_BUILDER, lws_now_secs(), 0))
goto bail;
- sul = malloc(sizeof(*sul));
- if (!sul) {
+ if (sais_add_to_inflight_list_if_absent(vhd, cb, task_template->uuid)) {
sais_task_reset(vhd, task_template->uuid, 1);
goto bail;
}
- memset(sul, 0, sizeof(*sul));
- lws_strncpy(sul->uuid, task_template->uuid, sizeof(sul->uuid));
- lws_dll2_add_tail(&sul->list, &cb->inflight_owner);
/* provisionally decrement until we hear from builder */
if (cb->avail_slots > 0)
@@ -1089,6 +1174,7 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
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));
+
sais_list_builders(vhd);
/* advance the task state first time we get logs */
@@ -1197,16 +1283,21 @@ sais_activity_cb(lws_sorted_usec_list_t *sul)
int
sais_continue_task(struct vhd *vhd, const char *task_uuid)
{
- char event_uuid[33], esc_uuid[129], *p, *start, url[128], mirror_path[256],
- update[128];
- lws_dll2_owner_t o, o_event;
+ char event_uuid[33], esc_uuid[129], *p, *start, url[128], mirror_path[256], update[128];
sai_task_t *task = NULL, *task_template;
- sai_event_t *event;
- sai_plat_t *cb;
- struct pss *pss;
+ lws_dll2_owner_t o, o_event;
+ struct lwsac *ac = NULL;
+ sai_uuid_list_t *ul;
sqlite3 *pdb = NULL;
+ sai_event_t *event;
int n, build_step;
- struct lwsac *ac = NULL;
+ struct pss *pss;
+ sai_plat_t *cb;
+
+ if (sais_is_task_inflight(vhd, task_uuid, &ul) && ul->started) {
+ lwsl_notice("%s: not continuing %s as listed on inflight\n", __func__, task_uuid);
+ return 1;
+ }
memset(event_uuid, 0, sizeof(event_uuid));
sai_task_uuid_to_event_uuid(event_uuid, task_uuid);
@@ -1214,9 +1305,6 @@ sais_continue_task(struct vhd *vhd, const char *task_uuid)
if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb) || !pdb)
return -1;
- sqlite3_exec(pdb, "ALTER TABLE tasks ADD COLUMN build_step INTEGER;",
- NULL, NULL, NULL);
-
lwsl_notice("%s: task_uuid %s, pdb %p\n", __func__, task_uuid, pdb);
lws_sql_purify(esc_uuid, task_uuid, sizeof(esc_uuid));
@@ -1351,6 +1439,14 @@ sais_continue_task(struct vhd *vhd, const char *task_uuid)
task->server_name = pss->server_name;
+ if (sais_add_to_inflight_list_if_absent(vhd, cb, task->uuid)) {
+ sais_task_reset(vhd, task->uuid, 1);
+ sais_event_db_close(vhd, &pdb);
+ lwsac_free(&task->ac_task_container);
+ free(task);
+ return -1;
+ }
+
lws_dll2_add_tail(&task->pending_assign_list, &pss->issue_task_owner);
lws_callback_on_writable(pss->wsi);
diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c
index e71cdf6..57207b2 100644
--- a/src/server/s-websrv.c
+++ b/src/server/s-websrv.c
@@ -238,7 +238,7 @@ sais_list_builders(struct vhd *vhd)
return 1;
}
- lwsl_warn("%s: count deserialized %d\n", __func__, (int)db_builders_owner.count);
+ // lwsl_warn("%s: count deserialized %d\n", __func__, (int)db_builders_owner.count);
p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p),
"{\"schema\":\"com.warmcat.sai.builders\",\"builders\":[");
@@ -255,8 +255,8 @@ sais_list_builders(struct vhd *vhd)
live_builder = sais_builder_from_uuid(vhd, builder_from_db->name, __FILE__, __LINE__);
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);
+ // 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;
lws_strncpy(builder_from_db->peer_ip, live_builder->peer_ip,
sizeof(builder_from_db->peer_ip));
@@ -264,10 +264,10 @@ sais_list_builders(struct vhd *vhd)
} else
builder_from_db->online = 0;
- if (builder_from_db->power_managed)
+ /* 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->stay_on); */
builder_from_db->powering_up = 0;
builder_from_db->powering_down = 0;
@@ -306,7 +306,7 @@ sais_list_builders(struct vhd *vhd)
p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}");
- lwsl_notice("%s: Broadcasting builder list: %s\n", __func__, vhd->json_builders);
+ // 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));
diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c
index 789a70b..3633e33 100644
--- a/src/server/s-ws-builder.c
+++ b/src/server/s-ws-builder.c
@@ -433,14 +433,16 @@ sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name)
int
sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl)
{
- char event_uuid[33], s[128], esc[96];
+ char event_uuid[33], s[128], esc[96], do_remove_uuid;
const sai_build_metric_t *metric;
sai_resource_requisition_t *rr;
sai_resource_wellknown_t *wk;
+ sai_plat_owner_t *bp_owner;
struct lwsac *ac = NULL;
sai_plat_t *build, *cb;
sai_rejection_t *rej;
sai_resource_t *res;
+ sai_uuid_list_t *ul;
lws_dll2_owner_t o;
sai_artifact_t *ap;
sai_task_t *task;
@@ -495,7 +497,7 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b
return 1;
}
- lwsl_hexdump_notice(buf, bl);
+ // lwsl_hexdump_notice(buf, bl);
if (m == LEJP_CONTINUE) {
pss->frag = 1;
@@ -516,10 +518,10 @@ handle:
* builder is sending us an array of platforms it provides us
*/
- pss->u.o = (sai_plat_owner_t *)pss->a.dest;
+ bp_owner = (sai_plat_owner_t *)pss->a.dest;
lws_start_foreach_dll(struct lws_dll2 *, pb,
- pss->u.o->plat_owner.head) {
+ bp_owner->plat_owner.head) {
build = lws_container_of(pb, sai_plat_t, sai_plat_list);
sai_plat_t *live_cb;
@@ -529,9 +531,9 @@ handle:
char q[1024];
lws_snprintf(q, sizeof(q),
- "INSERT INTO builders (name, platform, online, last_seen, peer_ip, sai_hash, lws_hash, windows) "
- "VALUES ('%s', '%s', 1, %llu, '%s', '%s', '%s', %d) "
- "ON CONFLICT(name) DO UPDATE SET online=1, last_seen=excluded.last_seen, "
+ "INSERT INTO builders (name, platform, last_seen, peer_ip, sai_hash, lws_hash, windows) "
+ "VALUES ('%s', '%s', %llu, '%s', '%s', '%s', %d) "
+ "ON CONFLICT(name) DO UPDATE SET last_seen=excluded.last_seen, "
"peer_ip=excluded.peer_ip, sai_hash=excluded.sai_hash, lws_hash=excluded.lws_hash",
build->name, build->platform, (unsigned long long)lws_now_secs(),
pss->peer_ip, build->sai_hash, build->lws_hash, build->windows);
@@ -550,20 +552,20 @@ 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));
lws_strncpy(live_cb->lws_hash, build->lws_hash,
sizeof(live_cb->lws_hash));
- live_cb->windows = build->windows;
- live_cb->online = 1;
- 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->s_inflight_count = (int)live_cb->inflight_owner.count;
- live_cb->s_last_rej_task_uuid[0] = '\0';
+ live_cb->windows = build->windows;
+ live_cb->online = 1;
+ live_cb->avail_slots = -1; /* ie, unknown */
+ 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->s_inflight_count = (int)live_cb->inflight_owner.count;
+ live_cb->s_last_rej_task_uuid[0] = '\0';
} else {
/* New builder, create a deep-copied, malloc'd object */
size_t nlen = strlen(build->name) + 1;
@@ -574,22 +576,23 @@ handle:
live_cb = malloc(sizeof(*live_cb) + nlen + plen);
if (live_cb) {
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);
}
@@ -619,7 +622,7 @@ handle:
goto bail;
}
} lws_end_foreach_dll(p);
-
+#if 0
lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) {
cb = lws_container_of(p, sai_plat_t, sai_plat_list);
if (cb->wsi == pss->wsi) {
@@ -628,6 +631,7 @@ handle:
goto bail;
}
} lws_end_foreach_dll(p);
+#endif
/*
* If we did allocate a task in pss->a.ac, responsibility of
@@ -650,6 +654,7 @@ bail:
log = (sai_log_t *)pss->a.dest;
sais_log_to_db(vhd, log);
+#if 0
if (pss->mark_started) {
pss->mark_started = 0;
pss->first_log_timestamp = log->timestamp;
@@ -657,6 +662,7 @@ bail:
SAIES_BEING_BUILT, 0, 0))
goto bail;
}
+#endif
if (log->finished) {
sai_plat_t *cb;
@@ -692,8 +698,7 @@ bail:
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)) {
- lws_dll2_remove(&sul->list);
- free(sul);
+ sais_inflight_entry_destroy(sul);
break;
}
} lws_end_foreach_dll_safe(d, d1);
@@ -740,12 +745,14 @@ bail:
case SAIM_WSSCH_BUILDER_TASKREJ:
/*
- * builder is updating us about his status, and may be
- * rejecting a task we tried to give him
+ * builder is updating us about a task status
*/
rej = (sai_rejection_t *)pss->a.dest;
+ if (!rej->task_uuid[0])
+ break;
+
rej->host_platform[sizeof(rej->host_platform) - 1] = '\0';
cb = sais_builder_from_uuid(vhd, rej->host_platform, __FILE__, __LINE__);
if (!cb) {
@@ -755,15 +762,47 @@ bail:
break;
}
- cb->avail_slots = rej->avail_slots;
- cb->avail_mem_kib = rej->avail_mem_kib;
- cb->avail_sto_kib = rej->avail_sto_kib;
+ cb->avail_slots = rej->avail_slots;
+ cb->avail_mem_kib = rej->avail_mem_kib;
+ cb->avail_sto_kib = rej->avail_sto_kib;
- lwsl_notice("%s: builder %s reports rejection (rej %s), slots %d, mem %d, sto %d\n",
- __func__, cb->name,
- rej->task_uuid[0] ? rej->task_uuid : "none",
+ lwsl_notice("%s: builder %s reports task status update, reason: %d, %s, slots %d, mem %d, sto %d\n",
+ __func__, cb->name, rej->reason, rej->task_uuid,
cb->avail_slots, cb->avail_mem_kib, cb->avail_sto_kib);
+ do_remove_uuid = 0;
+
+ switch (rej->reason) {
+ case SAI_TASK_REASON_ACCEPTED:
+ lwsl_notice("%s: SAI_TASK_REASON_ACCEPTED\n", __func__);
+ pss->first_log_timestamp = lws_now_secs();
+ 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;
+ break;
+ case SAI_TASK_REASON_BUSY:
+ lwsl_notice("%s: SAI_TASK_REASON_BUSY\n", __func__);
+ do_remove_uuid = 1;
+ break;
+ case SAI_TASK_REASON_DESTROYED:
+ lwsl_notice("%s: SAI_TASK_REASON_DESTROYED\n", __func__);
+ do_remove_uuid = 1;
+ break;
+ }
+
+ if (do_remove_uuid && sais_is_task_inflight(vhd, rej->task_uuid, &ul)) {
+ lwsl_notice("%s: ### Removing %s from inflight\n", __func__, rej->task_uuid);
+ lws_dll2_remove(&ul->list);
+ free(ul);
+ }
+
+
+#if 0
if (rej->task_uuid[0]) {
sai_uuid_list_t *sul;
@@ -781,11 +820,13 @@ bail:
sizeof(cb->last_rej_task_uuid));
sais_task_reset(vhd, rej->task_uuid, 1);
}
+#endif
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));
+
sais_list_builders(vhd);
lwsac_free(&pss->a.ac);
@@ -1329,18 +1370,14 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf,
* (all in .ac)
*/
- task = lws_container_of(pss->issue_task_owner.head, sai_task_t,
- pending_assign_list);
+ task = lws_container_of(pss->issue_task_owner.head, sai_task_t, pending_assign_list);
lws_dll2_remove(&task->pending_assign_list);
js = lws_struct_json_serialize_create(lsm_schema_map_ta,
LWS_ARRAY_SIZE(lsm_schema_map_ta),
0, task);
- if (!js) {
- lwsac_free(&task->ac_task_container);
- free(task);
- return 1;
- }
+ if (!js)
+ goto bail;
n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w);
lws_struct_json_serialize_destroy(&js);
@@ -1348,6 +1385,8 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf,
lwsac_free(&task->ac_task_container);
free(task);
+ lwsl_err("%s: ########## ATTACH TASK: %.*s\n", __func__, (int)w, start);
+
first = 1;
send_json:
@@ -1376,4 +1415,11 @@ send_json:
lws_callback_on_writable(pss->wsi);
return 0;
+
+bail:
+ lwsac_free(&task->ac_task_container);
+ free(task);
+
+ return 1;
+
}
diff --git a/src/web/CMakeLists.txt b/src/web/CMakeLists.txt
index 5035616..8f0ea81 100644
--- a/src/web/CMakeLists.txt
+++ b/src/web/CMakeLists.txt
@@ -100,6 +100,7 @@ if (requirements)
../../assets/arch-riscv64.svg
../../assets/arch-risc-v.svg
../../assets/arch-x86_64-amd.svg
+ ../../assets/arch-x86_64.svg
../../assets/x86_64.svg
../../assets/arch-x86_64-intel-i3.svg
../../assets/arch-x86_64-intel.svg
diff --git a/src/web/w-private.h b/src/web/w-private.h
index 3bcbc5a..3a8a7a7 100644
--- a/src/web/w-private.h
+++ b/src/web/w-private.h
@@ -151,8 +151,6 @@ struct pss {
sqlite3 *pdb_artifact;
sqlite3_blob *blob_artifact;
- lws_dll2_owner_t platform_owner; /* sai_platform_t builder offers */
- lws_dll2_owner_t task_cancel_owner; /* sai_platform_t builder offers */
lws_dll2_owner_t logs_owner;
lws_struct_args_t a;
diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c
index 87fd60a..fa125b8 100644
--- a/src/web/w-websrv.c
+++ b/src/web/w-websrv.c
@@ -115,8 +115,8 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags)
sai_browse_rx_evinfo_t *ei;
int n;
- lwsl_warn("%s: len %d, flags %d\n", __func__, (int)len, flags);
- lwsl_hexdump_notice(buf, len);
+ // lwsl_warn("%s: len %d, flags %d\n", __func__, (int)len, flags);
+ // lwsl_hexdump_notice(buf, len);
if (flags & LWSSS_FLAG_SOM) {
/* First fragment of a new message. Clear old parse results and init. */
@@ -298,8 +298,8 @@ saiw_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len,
*len = used;
*flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0);
- lwsl_ss_notice(m->ss, "Sending %d web->srv: ssflags %d", (int)*len, (int)*flags);
- lwsl_hexdump_notice(buf, *len);
+ // lwsl_ss_notice(m->ss, "Sending %d web->srv: ssflags %d", (int)*len, (int)*flags);
+ // lwsl_hexdump_notice(buf, *len);
if (m->wbltx)
return lws_ss_request_tx(m->ss);