Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / src / power / p-ws-server.c
Author[]Andy Green <andy@warmcat.com> 2025-10-09 12:49 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-10-09 18:20 UTC
Tree2a7132d0807d5de903dec7a8c31124aaf92277cf   Raw Patch
 
feat: Improve sai-server task allocation logic
feat: Improve sai-server task allocation logic
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;
Page fetched 0s ago, creation time: 9ms (vhost etag hits: 0%, cache hits: 0%)