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 / virt / v-private.h
Author[]Andy Green <andy@warmcat.com> 2025-10-11 04:38 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-10-12 05:34 UTC
Tree9624a0a4ebcb82274b7a5c5c49c0155945394dc2   Raw Patch
 
builder: refactor task tracking
builder: refactor task tracking
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);
Page fetched 0s ago, creation time: 9ms (vhost etag hits: 0%, cache hits: 0%)