diff --git a/assets/sai.js b/assets/sai.js
index 95d81fa..50db89f 100644
--- a/assets/sai.js
+++ b/assets/sai.js
@@ -1055,7 +1055,10 @@ function createBuilderDiv(plat) {
`<div class="res-bar"><div class="res-bar-inner res-bar-ram w-0"></div></div>` +
`<div class="res-bar"><div class="res-bar-inner res-bar-disk w-0"></div></div>` +
`</div>`;
- innerHTML += `<br>${plat.peer_ip}</td></tr></tbody></table>`;
+ innerHTML += `<br>${plat.peer_ip}<div class="server-state">` +
+ `Slots: ${plat.s_avail_slots}, In-flight: ${plat.s_inflight_count}<br>` +
+ `Last Reject: ${plat.s_last_rej_task_uuid ? plat.s_last_rej_task_uuid.substring(0, 8) : 'none'}` +
+ `</div></td></tr></tbody></table>`;
platDiv.innerHTML = innerHTML;
diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c
index cd3542a..591d02c 100644
--- a/src/builder/b-comms.c
+++ b/src/builder/b-comms.c
@@ -375,9 +375,17 @@ send_logs:
ns->task->uuid, (unsigned long long)lws_now_usecs(),
chunk->stdfd, (int)chunk->len);
- if (ns->finished_when_logs_drained && !ns->chunk_cache.count)
+ if (ns->finished_when_logs_drained && !ns->chunk_cache.count) {
n += lws_snprintf((char *)p + n, lws_ptr_diff_size_t(end, p) - (unsigned int)n,
"\"finished\":%d,", ns->retcode);
+ sp = ns->sp;
+ n += lws_snprintf((char *)p + n, lws_ptr_diff_size_t(end, p) - (unsigned int)n,
+ "\"avail_slots\":%d,\"avail_mem_kib\":%u,\"avail_sto_kib\":%u,",
+ (int)(sp->job_limit ? sp->job_limit : 6u) -
+ ((int)sp->nspawn_owner.count - 1),
+ saib_get_free_ram_kib(),
+ saib_get_free_disk_kib(builder.home));
+ }
n += lws_snprintf((char *)p + n, lws_ptr_diff_size_t(end, p) - (unsigned int)n,
"\"log\":\"");
diff --git a/src/builder/b-task.c b/src/builder/b-task.c
index 813a50d..8a7f5f1 100644
--- a/src/builder/b-task.c
+++ b/src/builder/b-task.c
@@ -313,6 +313,11 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm,
lws_snprintf(rej->host_platform, sizeof(rej->host_platform), "%s",
sp->name);
+ 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);
+
lws_dll2_add_tail(&rej->list, &spm->rejection_list);
return lws_ss_request_tx(spm->ss) ? -1 : 0;
diff --git a/src/common/include/private.h b/src/common/include/private.h
index 2f08956..f2b6e9d 100644
--- a/src/common/include/private.h
+++ b/src/common/include/private.h
@@ -238,6 +238,9 @@ 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;
} sai_rejection_t;
/*
@@ -293,6 +296,11 @@ typedef struct {
int finished;
int channel;
int uid;
+
+ /* builder can report this along with step completion */
+ int avail_slots;
+ unsigned int avail_mem_kib;
+ unsigned int avail_sto_kib;
} sai_log_t;
typedef struct {
@@ -422,6 +430,12 @@ typedef struct sai_plat_server_ref {
char was_active;
} sai_plat_server_ref_t;
+/* common struct for lists of task uuids on a builder */
+typedef struct sai_uuid_list {
+ lws_dll2_t list;
+ char uuid[65];
+} sai_uuid_list_t;
+
/*
* One of these instantiated per platform instance
*
@@ -460,6 +474,18 @@ typedef struct sai_plat {
int powering_down;
unsigned int job_limit;
+ /* server side only: builder resource tracking */
+ lws_dll2_owner_t inflight_owner; /* sai_uuid_list_t */
+ char last_rej_task_uuid[65];
+ int avail_slots;
+ unsigned int avail_mem_kib;
+ unsigned int avail_sto_kib;
+
+ /* server side only: for UI visibility */
+ int s_avail_slots;
+ int s_inflight_count;
+ char s_last_rej_task_uuid[65];
+
char windows;
char power_managed;
char stay_on;
@@ -585,10 +611,10 @@ extern const lws_struct_map_t
lsm_schema_map_plat_simple[1],
lsm_event[11],
lsm_task[29],
- lsm_log[7],
+ lsm_log[10],
lsm_artifact[8],
lsm_plat_list[1],
- lsm_task_rej[2],
+ lsm_task_rej[5],
lsm_task_cancel[1],
lsm_schema_json_map_can[1],
lsm_schema_json_map_task[1],
@@ -603,7 +629,7 @@ extern const lws_struct_map_t
;
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[13];
+extern const lws_struct_map_t lsm_plat_for_json[16];
extern const lws_ss_info_t ssi_said_logproxy;
extern struct lws_ss_handle *ssh[3];
diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c
index 0579b61..52b3bb0 100644
--- a/src/common/struct-metadata.c
+++ b/src/common/struct-metadata.c
@@ -111,6 +111,9 @@ const lws_struct_map_t lsm_plat_for_json[] = {
LSM_UNSIGNED(sai_plat_t, windows, "windows"),
LSM_UNSIGNED(sai_plat_t, power_managed, "power_managed"),
LSM_UNSIGNED(sai_plat_t, stay_on, "stay_on"),
+ LSM_SIGNED(sai_plat_t, s_avail_slots, "s_avail_slots"),
+ LSM_SIGNED(sai_plat_t, s_inflight_count, "s_inflight_count"),
+ LSM_CARRAY(sai_plat_t, s_last_rej_task_uuid, "s_last_rej_task_uuid"),
};
const lws_struct_map_t lsm_schema_map_plat_simple[] = {
@@ -208,6 +211,9 @@ const lws_struct_map_t lsm_schema_sq3_map_task[] = {
const lws_struct_map_t lsm_task_rej[] = {
LSM_CARRAY (sai_rejection_t, host_platform, "host_platform"),
LSM_CARRAY (sai_rejection_t, task_uuid, "task_uuid"),
+ 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"),
};
const lws_struct_map_t lsm_schema_json_task_rej[] = {
@@ -253,6 +259,9 @@ const lws_struct_map_t lsm_log[] = {
LSM_UNSIGNED (sai_log_t, finished, "finished"),
LSM_CARRAY (sai_log_t, task_uuid, "task_uuid"),
LSM_STRING_PTR (sai_log_t, log, "log"),
+ LSM_JO_SIGNED (sai_log_t, avail_slots, "avail_slots"),
+ LSM_JO_UNSIGNED (sai_log_t, avail_mem_kib, "avail_mem_kib"),
+ LSM_JO_UNSIGNED (sai_log_t, avail_sto_kib, "avail_sto_kib"),
};
const lws_struct_map_t lsm_schema_json_map_log[] = {
diff --git a/src/server/s-comms.c b/src/server/s-comms.c
index fa0e8c4..81cecfd 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 sai-web conn ###############");
/* remove pss from vhd->builders (active connection list) */
lws_dll2_remove(&pss->same);
diff --git a/src/server/s-private.h b/src/server/s-private.h
index 7707193..7b9e39f 100644
--- a/src/server/s-private.h
+++ b/src/server/s-private.h
@@ -348,3 +348,7 @@ sais_set_builder_power_state(struct vhd *vhd, const char *name, int up, int down
void
sais_mark_all_builders_offline(struct vhd *vhd);
+
+int
+sql3_get_string_cb(void *user, int cols, char **values, char **name);
+
diff --git a/src/server/s-task.c b/src/server/s-task.c
index 88dfada..5b4037b 100644
--- a/src/server/s-task.c
+++ b/src/server/s-task.c
@@ -332,7 +332,8 @@ sais_event_ran_platform(struct vhd *vhd, const char *event_uuid,
* Find the most recent task that still needs doing for platform, on any event
*/
static const sai_task_t *
-sais_task_pending(struct vhd *vhd, struct pss *pss, const char *platform)
+sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
+ const char *platform)
{
struct lwsac *ac = NULL, *failed_ac = NULL;
char esc_plat[96], pf[2048], query[384];
@@ -488,6 +489,14 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, const char *platform)
lws_sql_purify(esc_taskname, fti->taskname, sizeof(esc_taskname));
lws_snprintf(pf, sizeof(pf), " and (state == 0) and (platform == '%s') and (taskname == '%s')",
esc_plat, esc_taskname);
+ if (cb->last_rej_task_uuid[0]) {
+ char esc_uuid[130];
+
+ lws_sql_purify(esc_uuid, cb->last_rej_task_uuid,
+ sizeof(esc_uuid));
+ lws_snprintf(pf + strlen(pf), sizeof(pf) - strlen(pf),
+ " and (uuid != '%s')", esc_uuid);
+ }
lwsac_free(&pss->ac_alloc_task);
lws_dll2_owner_clear(&owner);
@@ -514,6 +523,14 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, const char *platform)
/* We have fallen back to doing tasks earliest-first */
lws_snprintf(pf, sizeof(pf), " and (state = 0) and (platform = '%s')", esc_plat);
+ if (cb->last_rej_task_uuid[0]) {
+ char esc_uuid[130];
+
+ lws_sql_purify(esc_uuid, cb->last_rej_task_uuid,
+ sizeof(esc_uuid));
+ lws_snprintf(pf + strlen(pf), sizeof(pf) - strlen(pf),
+ " and (uuid != '%s')", esc_uuid);
+ }
lwsac_free(&pss->ac_alloc_task);
lws_dll2_owner_t owner;
lws_dll2_owner_clear(&owner);
@@ -998,40 +1015,103 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb,
const char *platform_name)
{
const sai_task_t *task_template;
- sai_task_t *task = NULL;
+ char original_rejected_uuid[65];
+ sai_task_t temp_task;
+ sai_uuid_list_t *sul;
+ int attempts = 0;
+
+ if (cb->avail_slots <= 0) {
+ lwsl_info("%s: builder %s has no available slots\n", __func__,
+ cb->name);
+ return 1;
+ }
- /*
- * Look for a task for this platform, on any event that needs building
- */
+ lws_strncpy(original_rejected_uuid, cb->last_rej_task_uuid,
+ sizeof(original_rejected_uuid));
- task_template = sais_task_pending(vhd, pss, platform_name);
- if (!task_template)
- return 1;
+ while (attempts++ < 10) {
- lwsl_notice("%s: %s: task found %s\n", __func__, platform_name, cb->name);
+ /*
+ * Look for a task for this platform, on any event that needs building
+ */
- /* yes, we will offer it to him */
+ task_template = sais_task_pending(vhd, pss, cb, platform_name);
+ if (!task_template) {
+ lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid,
+ sizeof(cb->last_rej_task_uuid));
+ return 1;
+ }
- if (sais_set_task_state(vhd, cb->name, cb->name, task_template->uuid,
- SAIES_PASSED_TO_BUILDER, lws_now_secs(), 0))
- goto bail;
+ /*
+ * We have a candidate task, check if the builder has enough
+ * resources for it
+ */
+ memcpy(&temp_task, task_template, sizeof(temp_task));
+ sais_get_task_metrics_estimates(vhd, &temp_task);
+
+ if (temp_task.est_peak_mem_kib > cb->avail_mem_kib ||
+ temp_task.est_disk_kib > cb->avail_sto_kib) {
+ lwsl_notice("%s: builder %s lacks resources for task %s "
+ "(mem %uk/%uk, sto %uk/%uk), trying another\n",
+ __func__, cb->name, temp_task.uuid,
+ temp_task.est_peak_mem_kib, cb->avail_mem_kib,
+ temp_task.est_disk_kib, cb->avail_sto_kib);
+
+ /* mark it rejected for this builder and try again */
+ lws_strncpy(cb->last_rej_task_uuid, temp_task.uuid,
+ sizeof(cb->last_rej_task_uuid));
+ continue;
+ }
- /* advance the task state first time we get logs */
- pss->mark_started = 1;
+ lwsl_notice("%s: %s: task %s found for %s\n", __func__,
+ platform_name, task_template->uuid, cb->name);
- sais_continue_task(vhd, task_template->uuid);
+ /* yes, we will offer it to him */
- /*
- * We are going to leave here with a live pss->a.ac (pointed into by
- * task->one_event) that the caller has to take responsibility to
- * clean up pss->a.ac
- */
+ if (sais_set_task_state(vhd, cb->name, cb->name, task_template->uuid,
+ SAIES_PASSED_TO_BUILDER, lws_now_secs(), 0))
+ goto bail;
- return 0;
+ sul = malloc(sizeof(*sul));
+ if (!sul) {
+ 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)
+ cb->avail_slots--;
+
+ 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);
+
+ /* advance the task state first time we get logs */
+ pss->mark_started = 1;
+
+ sais_continue_task(vhd, task_template->uuid);
+
+ lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid,
+ sizeof(cb->last_rej_task_uuid));
+
+ return 0;
+ }
+
+ lwsl_warn("%s: exceeded max attempts to find suitable task for %s\n",
+ __func__, cb->name);
+ lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid,
+ sizeof(cb->last_rej_task_uuid));
+
+ return 1;
bail:
- if (task)
- free(task);
+ lws_strncpy(cb->last_rej_task_uuid, original_rejected_uuid,
+ sizeof(cb->last_rej_task_uuid));
lwsac_free(&pss->a.ac);
lwsac_free(&pss->ac_alloc_task);
diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c
index eca88e9..789a70b 100644
--- a/src/server/s-ws-builder.c
+++ b/src/server/s-ws-builder.c
@@ -495,6 +495,8 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b
return 1;
}
+ lwsl_hexdump_notice(buf, bl);
+
if (m == LEJP_CONTINUE) {
pss->frag = 1;
return 0;
@@ -556,6 +558,12 @@ handle:
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';
} else {
/* New builder, create a deep-copied, malloc'd object */
size_t nlen = strlen(build->name) + 1;
@@ -576,6 +584,10 @@ handle:
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;
lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip));
@@ -647,6 +659,54 @@ bail:
}
if (log->finished) {
+ sai_plat_t *cb;
+ 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';
+ lws_sql_purify(esc_uuid, log->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]) {
+ cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__);
+ if (cb) {
+ 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';
+
+ 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);
+ break;
+ }
+ } lws_end_foreach_dll_safe(d, d1);
+
+ 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);
+ }
+ }
+ sais_event_db_close(vhd, &pdb);
+ }
/*
* We have reached the end of the logs for this task
*/
@@ -695,12 +755,38 @@ bail:
break;
}
- lwsl_notice("%s: builder %s reports rejection (rej %s)\n",
+ 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");
+ rej->task_uuid[0] ? rej->task_uuid : "none",
+ cb->avail_slots, cb->avail_mem_kib, cb->avail_sto_kib);
+
+ if (rej->task_uuid[0]) {
+ sai_uuid_list_t *sul;
+
+ 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, rej->task_uuid)) {
+ lws_dll2_remove(&sul->list);
+ free(sul);
+ break;
+ }
+ } lws_end_foreach_dll_safe(d, d1);
- if (rej->task_uuid[0])
+ lws_strncpy(cb->last_rej_task_uuid, rej->task_uuid,
+ sizeof(cb->last_rej_task_uuid));
sais_task_reset(vhd, rej->task_uuid, 1);
+ }
+
+ 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);
break;