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
Author[]google-labs-jules[bot] <161369871+google-labs-jules[bot]@users.noreply.github.c...> 2025-09-01 15:40 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-09-02 15:10 UTC
Tree51ef0884c6997da65ecec824a43ddf3649a9249f   Raw Patch
 
fix(builder): Reject task if already running
fix(builder): Reject task if already running
diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index 62ccd72..19c45f2 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -523,8 +523,8 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, const char *pq; int n; - lwsl_user("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name(state), - (unsigned int)ack); + // lwsl_user("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name(state), + // (unsigned int)ack); switch (state) { diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index 6935808..5d47173 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -308,6 +308,12 @@ skip: if (op->spawn) free(op->spawn); + if (ns->spm) { + ns->spm->phase = PHASE_START_ATTACH; + if (lws_ss_request_tx(ns->spm->ss)) + lwsl_warn("%s: lws_ss_request_tx failed\n", __func__); + } + /* * add a final zero-length log with the retcode to the list of pending * logs diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index 28a478e..aabdb46 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -307,8 +307,8 @@ saib_power_stay_rx(void *userobj, const uint8_t *buf, size_t len, int flags) */ lws_sul_cancel(&builder.sul_idle); - lwsl_warn("%s: %s: stay applied: cancelled idle grace time\n", - __func__, builder.host); + // lwsl_warn("%s: %s: stay applied: cancelled idle grace time\n", + // __func__, builder.host); } else { /* diff --git a/src/builder/b-task.c b/src/builder/b-task.c index 941b291..0fb8c71 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -157,8 +157,8 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, */ if (rej_task_uuid) { - lwsl_notice("%s: builder .%d occupied reject\n", - __func__, (int)(intptr_t)spm->opaque_data); + lwsl_notice("%s: builder %s occupied reject\n", + __func__, sp->name); lws_strncpy(rej->task_uuid, rej_task_uuid, sizeof(rej->task_uuid)); @@ -266,17 +266,22 @@ saib_sub_cleaner_cb(lws_sorted_usec_list_t *sul) { struct sai_nspawn *ns = lws_container_of(sul, struct sai_nspawn, sul_cleaner); - lwsl_warn("%s: Task completion grace period ended with ns alive\n", __func__); + lwsl_warn("%s: +++++ Task completion grace period ended with ns alive\n", __func__); - // saib_task_destroy(ns); - - if (ns->op && ns->op->lsp) + + if (ns->op && ns->op->lsp) { + lwsl_notice("%s: +++++++++++ killing child process\n", __func__); lws_spawn_piped_kill_child_process(ns->op->lsp); + } else { + lwsl_err("%s: ========================== unable to kill child process (already dead?)\n", __func__); + saib_task_destroy(ns); + } } void saib_task_grace(struct sai_nspawn *ns) { + lwsl_err("%s: +++++ starting task grace wait\n", __func__); ns->finished_when_logs_drained = 1; /* destroy ns if logs are gone + spawn reaped */ lws_sul_schedule(builder.context, 0, &ns->sul_cleaner, saib_sub_cleaner_cb, 20 * LWS_USEC_PER_SEC); @@ -583,10 +588,12 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) 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); - if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { - lwsl_err("%s: server offered task that's already running\n", __func__); + 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); @@ -612,8 +619,9 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) 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); if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { - ns = xns; - break; + lwsl_warn("%s: server offered task that's already running\n", __func__); + saib_queue_task_status_update(sp, spm, task->uuid); + return 0; } if (!xns->task && !ns) ns = xns; @@ -735,8 +743,8 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) lwsac_free(&ns->task->ac_task_container); ns->task = task; /* we are owning this nspawn for the duration */ - if (!task->build_step) { - ns->current_step = 0; + ns->current_step = task->build_step; + if (!ns->current_step) { ns->spins = 0; ns->user_cancel = 0; ns->us_cpu_user = 0; @@ -840,7 +848,6 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) can = (sai_cancel_t *)a.dest; lwsl_notice("%s: received task cancel for %s\n", __func__, can->task_uuid); - lwsl_notice("%s: SAIB_RX_TASK_CANCEL: %s\n", __func__, can->task_uuid); lws_start_foreach_dll_safe(struct lws_dll2 *, mp, mp1, builder.sai_plat_owner.head) { diff --git a/src/common/include/private.h b/src/common/include/private.h index 971ac19..6087e96 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -557,4 +557,3 @@ sul_idle_cb(lws_sorted_usec_list_t *sul); int sai_uuid16_create(struct lws_context *context, char *dest33); - diff --git a/src/server/s-central.c b/src/server/s-central.c index 9ccb5c6..5f30bbb 100644 --- a/src/server/s-central.c +++ b/src/server/s-central.c @@ -132,7 +132,7 @@ sais_central_clean_abandoned(struct vhd *vhd) if (task_uuid) { lwsl_notice("%s: resetting abandoned task %s\n", __func__, (const char *)task_uuid); - sais_task_reset(vhd, (const char *)task_uuid); + sais_task_reset(vhd, (const char *)task_uuid, 0); } } sqlite3_finalize(sm); diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 573c3e4..e122be3 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -58,10 +58,6 @@ typedef struct sai_job { } sai_job_t; -const lws_struct_map_t lsm_schema_map_ta[] = { - LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), -}; - extern const lws_struct_map_t lsm_schema_sq3_map_event[]; extern const lws_ss_info_t ssi_server; diff --git a/src/server/s-private.h b/src/server/s-private.h index 44df052..5387e95 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -288,7 +288,7 @@ void sais_activity_cb(lws_sorted_usec_list_t *sul); sai_db_result_t -sais_task_reset(struct vhd *vhd, const char *task_uuid); +sais_task_reset(struct vhd *vhd, const char *task_uuid, int from_rejection); int sais_task_cancel(struct vhd *vhd, const char *task_uuid); diff --git a/src/server/s-task.c b/src/server/s-task.c index 11712ef..383db95 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -198,6 +198,11 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, sais_taskchange(vhd->h_ss_websrv, task_uuid, state); + if (state == SAIES_SUCCESS || state == SAIES_FAIL || + state == SAIES_CANCELLED) + lws_sul_schedule(vhd->context, 0, &vhd->sul_central, + sais_central_cb, 1); + sais_platforms_with_tasks_pending(vhd); /* @@ -757,30 +762,70 @@ sais_task_cancel(struct vhd *vhd, const char *task_uuid) static int sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) { + char event_uuid[33], builder_name[128], esc_uuid[129], q[128]; + struct pss *pss_match = NULL; + sai_plat_t *cb; + sqlite3 *pdb = NULL; sai_cancel_t *can; lwsl_notice("%s: builders count %d\n", __func__, vhd->builders.count); /* - * For every pss that we have from builders... + * We will send the task cancel message only to the builder that was + * assigned the task, if any. */ - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->builders.head) { - struct pss *pss = lws_container_of(p, struct pss, same); + sai_task_uuid_to_event_uuid(event_uuid, task_uuid); + + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + lwsl_err("%s: unable to open event-specific database\n", __func__); + return -1; + } + + builder_name[0] = '\0'; + lws_sql_purify(esc_uuid, task_uuid, sizeof(esc_uuid)); + lws_snprintf(q, sizeof(q), "select builder_name from tasks where uuid='%s'", + esc_uuid); + if (sqlite3_exec(pdb, q, sql3_get_string_cb, builder_name, NULL) != + SQLITE_OK || + !builder_name[0]) { + sais_event_db_close(vhd, &pdb); /* - * ... queue the task cancel message + * This is not an error... the task may not have had a builder + * assigned yet. There's nothing to do. */ - can = malloc(sizeof *can); - if (!can) - return -1; - memset(can, 0, sizeof(*can)); + return 0; + } + sais_event_db_close(vhd, &pdb); - lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); - lws_dll2_add_tail(&can->list, &pss->task_cancel_owner); - lws_callback_on_writable(pss->wsi); + cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); + if (!cb) + /* Builder not connected, nothing to do */ + return 0; + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->builders.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + if (pss->wsi == cb->wsi) { + pss_match = pss; + break; + } } lws_end_foreach_dll(p); + if (!pss_match) + /* Builder is live but has no pss? */ + return 0; + + can = malloc(sizeof *can); + if (!can) + return -1; + + memset(can, 0, sizeof(*can)); + + lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); + + lws_dll2_add_tail(&can->list, &pss_match->task_cancel_owner); + lws_callback_on_writable(pss_match->wsi); + return 0; } @@ -790,7 +835,7 @@ sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) */ sai_db_result_t -sais_task_reset(struct vhd *vhd, const char *task_uuid) +sais_task_reset(struct vhd *vhd, const char *task_uuid, int from_rejection) { char esc[96], cmd[256], event_uuid[33]; sqlite3 *pdb = NULL; @@ -842,15 +887,18 @@ sais_task_reset(struct vhd *vhd, const char *task_uuid) sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 1, 1); - lwsl_notice("%s: stopping task %s on builders\n", __func__, task_uuid); - lwsl_notice("%s: stopping task %s on builders\n", __func__, task_uuid); sais_task_stop_on_builders(vhd, task_uuid); /* - * Reassess now if there's a builder we can match to a pending task + * Reassess now if there's a builder we can match to a pending task, + * but not if we are being reset due to a rejection... that would + * just cause us to spam the builder with the same task again */ - lws_sul_schedule(vhd->context, 0, &vhd->sul_central, sais_central_cb, 1); + if (!from_rejection) { + lwsl_err("%s: scheduling sul_central to find a new task\n", __func__); + lws_sul_schedule(vhd->context, 0, &vhd->sul_central, sais_central_cb, 1); + } /* * Recompute startable task platforms and broadcast to all sai-power, @@ -1107,6 +1155,7 @@ sais_continue_task(struct vhd *vhd, const char *task_uuid) } if (!p) { /* no more steps */ + lwsl_err("%s: +++++++++++++++++++ determined no more steps after build_step %d for task %s, setting SAIES_SUCCESS\n", __func__, build_step, task->uuid); sais_set_task_state(vhd, NULL, NULL, task->uuid, SAIES_SUCCESS, 0, 0); sais_event_db_close(vhd, &pdb); lwsac_free(&task->ac_task_container); diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index e96f50e..b09ac29 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -406,7 +406,7 @@ sais_event_reset(struct vhd *vhd, const char *event_uuid) lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { sai_task_t *t = lws_container_of(p, sai_task_t, list); - if (sais_task_reset(vhd, t->uuid) == SAI_DB_RESULT_BUSY) { + if (sais_task_reset(vhd, t->uuid, 0) == SAI_DB_RESULT_BUSY) { sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); sais_event_db_close(vhd, &pdb); lwsac_free(&ac); @@ -529,7 +529,7 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) lws_start_foreach_dll(struct lws_dll2 *, p, o.head) { sai_task_t *t = lws_container_of(p, sai_task_t, list); - if (sais_task_reset(vhd, t->uuid) == SAI_DB_RESULT_BUSY) { + if (sais_task_reset(vhd, t->uuid, 0) == SAI_DB_RESULT_BUSY) { sqlite3_exec(pdb, "END TRANSACTION", NULL, NULL, &err); sais_event_db_close(vhd, &pdb); lwsac_free(&ac); @@ -591,7 +591,7 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) goto soft_error; lwsl_ss_warn(m->ss, "SAIS_WS_WEBSRV_RX_TASKRESET: %s: received", ei->event_hash); - if (sais_task_reset(m->vhd, ei->event_hash)) + if (sais_task_reset(m->vhd, ei->event_hash, 0)) lwsl_ss_err(m->ss, "taskreset failed"); break; } diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 63211af..5e2d59b 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -30,20 +30,9 @@ #include "s-private.h" #include "s-metrics-db.h" -typedef struct { - int count; -} count_ctx_t; - -#if 0 -static int -online_builder_count_cb(void *priv, int cols, char **cv, char **cn) -{ - count_ctx_t *ctx = (count_ctx_t *)priv; - ctx->count++; - lwsl_err("%s: FOUND an online builder in DB: %s\n", __func__, cv[0]); - return 0; -} -#endif +const lws_struct_map_t lsm_schema_map_ta[] = { + LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), +}; enum sai_overview_state { SOS_EVENT, @@ -410,7 +399,7 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) if (task_uuid) { lwsl_notice("%s: resetting task %s from disconnected builder %s\n", __func__, (const char *)task_uuid, cb->name); - sais_task_reset(vhd, (const char *)task_uuid); + sais_task_reset(vhd, (const char *)task_uuid, 0); } } sqlite3_finalize(sm); @@ -723,8 +712,8 @@ bail: __func__, cb->name, rej->task_uuid[0] ? rej->task_uuid : "none"); - if (rej->task_uuid[0]) - sais_task_reset(vhd, rej->task_uuid); +// if (rej->task_uuid[0]) +// sais_task_reset(vhd, rej->task_uuid, 1); lwsac_free(&pss->a.ac); break; @@ -1134,7 +1123,6 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, sai_cancel_t *c = lws_container_of(pss->task_cancel_owner.head, sai_cancel_t, list); - lwsl_notice("%s: sending task cancel for %s\n", __func__, c->task_uuid); js = lws_struct_json_serialize_create(lsm_schema_json_map_can, LWS_ARRAY_SIZE(lsm_schema_json_map_can), 0, c); if (!js)
Page fetched 0s ago, creation time: 9ms (vhost etag hits: 0%, cache hits: 0%)