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 / assets / builder-instance.svg
Author[]Andy Green <andy@warmcat.com> 2025-11-06 12:04 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-11-06 14:18 UTC
Treeed19601ba845bd708209ea8aff67d3111233ded3   Raw Patch
 
clean: reduce indent levels
clean: reduce indent levels
diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 6ff64b2..ee3f039 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -118,7 +118,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, if (!lws_pvo_get_str(in, "task-abandoned-timeout-mins", &num)) vhd->task_abandoned_timeout_mins = (unsigned int)atoi(num); else - vhd->task_abandoned_timeout_mins = 3000; + vhd->task_abandoned_timeout_mins = 8 * 60; if (lws_pvo_get_str(in, "database", &vhd->sqlite3_path_lhs)) { lwsl_err("%s: database pvo required\n", __func__); @@ -274,7 +274,6 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, goto passthru; } - lwsl_user("LWS_CALLBACK_HTTP_BODY: %d\n", (int)len); /* create the POST argument parser if not already existing */ if (!pss->spa) { @@ -345,8 +344,6 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, break; case LWS_CALLBACK_HTTP_BODY_COMPLETION: - lwsl_user("%s: LWS_CALLBACK_HTTP_BODY_COMPLETION: %d\n", - __func__, (int)len); if (!pss->our_form) { lwsl_user("%s: no sai form\n", __func__); @@ -359,36 +356,29 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, pss->spa = NULL; } - if (pss->spa_failed) - lwsl_notice("%s: notification failed\n", __func__); - else { - lwsl_notice("%s: notification: %d %s %s %s\n", __func__, - pss->sn.action, pss->sn.e.hash, - pss->sn.e.ref, pss->sn.e.repo_name); + /* + * Inform sai-webs about notification processing, so + * they can update connected browsers to show the new + * event + */ + n = lws_snprintf((char *)start, sizeof(buf) - LWS_PRE, + "{\"schema\":\"sai-overview\"}"); - /* - * Inform sai-webs about notification processing, so - * they can update connected browsers to show the new - * event - */ - n = lws_snprintf((char *)start, sizeof(buf) - LWS_PRE, - "{\"schema\":\"sai-overview\"}"); + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = start; + info.len = (size_t)n; + info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - memset(&info, 0, sizeof(info)); - info.private_source_idx = SAI_WEBSRV_PB__GENERATED; - info.buf = start; - info.len = (size_t)n; - info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) - lwsl_warn("%s: buflist append failed\n", __func__); - } + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: buflist append failed\n", __func__); if (lws_return_http_status(wsi, - pss->spa_failed ? HTTP_STATUS_FORBIDDEN : - HTTP_STATUS_OK, - NULL) < 0) + pss->spa_failed ? HTTP_STATUS_FORBIDDEN : + HTTP_STATUS_OK, + NULL) < 0) return -1; + return 0; /* @@ -401,18 +391,17 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, case LWS_CALLBACK_ESTABLISHED: pss->wsi = wsi; pss->vhd = vhd; + if (!vhd) return -1; if (lws_hdr_total_length(wsi, WSI_TOKEN_GET_URI)) { - if (lws_hdr_copy(wsi, (char *)start, 64, - WSI_TOKEN_GET_URI) < 0) + if (lws_hdr_copy(wsi, (char *)start, 64, WSI_TOKEN_GET_URI) < 0) return -1; } #if defined(LWS_ROLE_H2) else - if (lws_hdr_copy(wsi, (char *)start, 64, - WSI_TOKEN_HTTP_COLON_PATH) < 0) + if (lws_hdr_copy(wsi, (char *)start, 64, WSI_TOKEN_HTTP_COLON_PATH) < 0) return -1; #endif @@ -462,7 +451,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 builder conn ###############"); + lwsl_wsi_user(wsi, "#### sai-server: CLOSED builder conn ####"); /* remove pss from vhd->builders (active connection list) */ lws_dll2_remove(&pss->same); @@ -484,7 +473,6 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, * Update the sai-webs about the builder removal, so they * can update their connected browsers */ - lwsl_wsi_warn(pss->wsi, "LWS_CALLBACK_CLOSED: doing WSS_PREPARE_BUILDER_SUMMARY\n"); sais_list_builders(vhd); break; @@ -517,12 +505,10 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, return -1; if (!pss->announced) { - /* * Update the sai-webs about the builder creation, so * they can update their connected browsers */ - lwsl_wsi_warn(pss->wsi, "LWS_CALLBACK_RECEIVE: unannounced pss doing WSS_PREPARE_BUILDER_SUMMARY\n"); sais_list_builders(vhd); pss->announced = 1; diff --git a/src/server/s-private.h b/src/server/s-private.h index a9e51e8..4b9d2c1 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -322,7 +322,7 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, const char *cns_name); int -sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char force); +sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid); int sais_set_task_state(struct vhd *vhd, const char *task_uuid, sai_event_state_t state, diff --git a/src/server/s-task-helpers.c b/src/server/s-task-helpers.c index f7b7895..e77d78b 100644 --- a/src/server/s-task-helpers.c +++ b/src/server/s-task-helpers.c @@ -347,7 +347,7 @@ sais_set_task_state(struct vhd *vhd, const char *task_uuid, if (ostate == SAIES_STEP_SUCCESS) { lwsl_notice("%s: sais_set_task_state() is calling sais_create_and_offer_task_step()\n", __func__); - sais_create_and_offer_task_step(vhd, task_uuid, 1); + sais_create_and_offer_task_step(vhd, task_uuid); } return 0; diff --git a/src/server/s-task.c b/src/server/s-task.c index b80b6cd..8b158ed 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -27,20 +27,14 @@ #include "s-private.h" -/* temporary info about a task that failed in a previous run */ -typedef struct sai_failed_task_info { - lws_dll2_t list; - /* over-allocated */ - const char *build; - const char *taskname; -} sai_failed_task_info_t; + /* * Checks if a given event db contains any tasks for a given platform */ static int -sais_event_ran_platform(struct vhd *vhd, const char *event_uuid, - const char *platform) +sais_event_check_for_plat_tasks(struct vhd *vhd, const char *event_uuid, + const char *platform) { sqlite3 *check_pdb = NULL; char query[256]; @@ -180,6 +174,12 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, struct lwsac *ac = NULL, *failed_ac = NULL; char esc_plat[96], pf[2048], query[384]; lws_dll2_owner_t o, failed_tasks_owner; + typedef struct sai_failed_task_info { + lws_dll2_t list; + /* over-allocated */ + const char *build; + const char *taskname; + } sai_failed_task_info_t; unsigned int pending_count; int n; @@ -221,200 +221,204 @@ sais_task_pending(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, uint64_t last_created; int m; - lwsl_notice("candidate event %s '%s'\n", e->uuid, esc_plat); + // lwsl_notice("candidate event %s '%s'\n", e->uuid, esc_plat); - if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { + if (sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) + goto next; + + /* + * 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 or state = 9) and platform = '%s'", esc_plat); + m = sqlite3_exec(pdb, query, sql3_get_integer_cb, &pending_count, NULL); + + if (m != SQLITE_OK) { + pending_count = 0; + lwsl_err("%s: query failed: %d\n", __func__, m); + } + + // lwsl_notice("%s: %s: platform: '%s' startable tasks: %d\n", __func__, e->uuid, esc_plat, pending_count); + + if (pending_count <= 0) { + lwsl_notice("%s: platform %s: no pending count\n", __func__, platform); + goto close_next; + } + + /* 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)); + last_created = e->created; + + do { + sqlite3_stmt *sm; + int pr; + + prev_event_uuid[0] = '\0'; + lws_snprintf(query, sizeof(query), + "select uuid, created from events where repo_name='%s' and " + "ref='%s' and created < %llu " + "order by created desc limit 1", + esc_repo, esc_ref, (unsigned long long)last_created); /* - * Find out how many tasks in startable state for this platform, - * on this event + * ... this is the 32-char EVENT uuid coming, + * not a compound (64 char) task one */ - lws_snprintf(query, sizeof(query), "select count(state) from tasks where " - "(state = 0 or state = 9) and platform = '%s'", esc_plat); - m = sqlite3_exec(pdb, query, sql3_get_integer_cb, &pending_count, NULL); - - if (m != SQLITE_OK) { - pending_count = 0; - lwsl_err("%s: query failed: %d\n", __func__, m); + pr = sqlite3_prepare_v2(vhd->server.pdb, query, -1, &sm, NULL); + if (pr != SQLITE_OK) { + lwsl_warn("%s: sq3 prep returned %d instead of SQLITE_OK\n", __func__, pr); + break; } + if (sqlite3_step(sm) == SQLITE_ROW) { + const char *u = (const char *)sqlite3_column_text(sm, 0); - lwsl_notice("%s: %s: platform: '%s' startable tasks: %d\n", __func__, e->uuid, esc_plat, pending_count); + if (u) + lws_strncpy(prev_event_uuid, (const char *)u, sizeof(prev_event_uuid)); - if (pending_count > 0) { /* there are some startable tasks on this event */ + last_created = (uint64_t)sqlite3_column_int64(sm, 1); + } else + lwsl_notice("%s: no results from event check %s %s\n", __func__, esc_repo, esc_ref); - lws_sql_purify(esc_repo, e->repo_name, sizeof(esc_repo)); - lws_sql_purify(esc_ref, e->ref, sizeof(esc_ref)); - last_created = e->created; - - do { - sqlite3_stmt *sm; - int pr; + sqlite3_finalize(sm); - prev_event_uuid[0] = '\0'; - lws_snprintf(query, sizeof(query), - "select uuid, created from events where repo_name='%s' and " - "ref='%s' and created < %llu " - "order by created desc limit 1", - esc_repo, esc_ref, (unsigned long long)last_created); + if (!prev_event_uuid[0]) { + lwsl_notice("%s: breaking due to NUL prev_event_uuid\n", __func__); + break; + } - /* this is the 32-char EVENT uuid coming, not a compound (64 char) task one */ + if (!sais_event_check_for_plat_tasks(vhd, prev_event_uuid, esc_plat)) { + lwsl_notice("%s: continuing due to event_ran_platform 0\n", __func__); + continue; + } - pr = sqlite3_prepare_v2(vhd->server.pdb, query, -1, &sm, NULL); - if (pr != SQLITE_OK) { - lwsl_warn("%s: sq3 prep returned %d instead of SQLITE_OK\n", __func__, pr); - break; - } - if (sqlite3_step(sm) == SQLITE_ROW) { - const char *u = (const char *)sqlite3_column_text(sm, 0); - if (u) { -#if 0 - if (sais_is_task_inflight(vhd, NULL, 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; - } -#endif - lws_strncpy(prev_event_uuid, (const char *)u, sizeof(prev_event_uuid)); - } - last_created = (uint64_t)sqlite3_column_int64(sm, 1); - } else - lwsl_notice("%s: no results from event check %s %s\n", __func__, esc_repo, esc_ref); - - sqlite3_finalize(sm); - - if (!prev_event_uuid[0]) { - lwsl_notice("%s: breaking due to NUL prev_event_uuid\n", __func__); - break; - } - if (!sais_event_ran_platform(vhd, prev_event_uuid, esc_plat)) { - lwsl_notice("%s: continuing due to event_ran_platform 0\n", __func__); - continue; - } + lws_strncpy(checked_uuid, prev_event_uuid, sizeof(checked_uuid)); + break; + } while (1); - lws_strncpy(checked_uuid, prev_event_uuid, sizeof(checked_uuid)); - break; - } while (1); - - if (checked_uuid[0]) { - lwsl_notice("%s: checked_uuid %s\n", __func__, checked_uuid); - if (!sais_event_db_ensure_open(vhd, checked_uuid, 1, &prev_pdb)) { - sqlite3_stmt *sm; - - /* we are looking for failed tasks here */ - - lws_snprintf(query, sizeof(query), - "select taskname from tasks where " - "state = 4 and platform = ?"); - - if (sqlite3_prepare_v2(prev_pdb, query, -1, &sm, NULL) == SQLITE_OK) { - const unsigned char *t; - sai_failed_task_info_t *fti; - - sqlite3_bind_text(sm, 1, esc_plat, -1, SQLITE_TRANSIENT); - - while (1) { - int nn = sqlite3_step(sm); - - if (nn != SQLITE_ROW) - break; - - t = sqlite3_column_text(sm, 0); - if (!t) - continue; - - /* - * We found errored tasks in the previous event for this - * repo / branch / platform. Let's record them in a temp - * lwsac and condsider if we should use this info to - * prioritize running the corresponding task in the current - * event first - */ - - fti = lwsac_use_zero(&failed_ac, sizeof(*fti) + - strlen((const char *)t) + 1, 256); - if (fti) { - fti->taskname = (const char *)&fti[1]; - memcpy((char *)fti->taskname, t, - strlen((const char *)t) + 1); - lws_dll2_add_tail(&fti->list, &failed_tasks_owner); - } - } - sqlite3_finalize(sm); - } else - lwsl_err("%s: query fail 1\n", __func__); - - sais_event_db_close(vhd, &prev_pdb); - } else - lwsl_err("%s: unable to open %s\n", __func__, checked_uuid); - } else - lwsl_notice("%s: platform %s: no checked uuid\n", __func__, platform); + if (checked_uuid[0] && + !sais_event_db_ensure_open(vhd, checked_uuid, 1, &prev_pdb)) { + sqlite3_stmt *sm; - /* - * Let's go through the tasks that failed last time we built this repo / branch, and see - * if we can find the analagous task in the current event. - */ + /* we are looking for failed tasks here */ - lws_start_foreach_dll(struct lws_dll2 *, p_fail, failed_tasks_owner.head) { - sai_failed_task_info_t *fti = lws_container_of(p_fail, sai_failed_task_info_t, list); - char esc_taskname[256]; - lws_dll2_owner_t owner; - - lws_sql_purify(esc_taskname, fti->taskname, sizeof(esc_taskname)); - lws_snprintf(pf, sizeof(pf), " and (state == 0 or state == 9) and (platform == '%s') and (taskname == '%s')", - esc_plat, esc_taskname); - - lwsac_free(&pss->ac_alloc_task); - lws_dll2_owner_clear(&owner); - n = lws_struct_sq3_deserialize(pdb, pf, NULL, - lsm_schema_sq3_map_task, - &owner, &pss->ac_alloc_task, 0, 1); - if (owner.count) { - lwsl_notice("%s: Prioritizing failed task for %s ('%s')\n", - __func__, platform, fti->taskname); - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - lwsac_free(&failed_ac); - memcpy(&pss->alloc_task, lws_container_of( - owner.head, sai_task_t, list), sizeof(pss->alloc_task)); - - lwsl_notice("%s: platform %s: returning selected task\n", __func__, platform); - - return &pss->alloc_task; - } - } lws_end_foreach_dll(p_fail); + lws_snprintf(query, sizeof(query), + "select taskname from tasks where " + "state = 4 and platform = ?"); - lwsl_notice("%s: no priority\n", __func__); + if (sqlite3_prepare_v2(prev_pdb, query, -1, &sm, NULL) == SQLITE_OK) { + const unsigned char *t; + sai_failed_task_info_t *fti; - /* We have fallen back to doing tasks earliest-first */ + sqlite3_bind_text(sm, 1, esc_plat, -1, SQLITE_TRANSIENT); - lws_snprintf(pf, sizeof(pf), " and (state = 0 or state = 9) and (platform = '%s')", esc_plat); + while (1) { + int nn = sqlite3_step(sm); - lwsac_free(&pss->ac_alloc_task); - lws_dll2_owner_t owner; - lws_dll2_owner_clear(&owner); - n = lws_struct_sq3_deserialize(pdb, pf, "uid asc ", - lsm_schema_sq3_map_task, - &owner, &pss->ac_alloc_task, 0, 1); - // lwsl_notice("%s: deser returned %d\n", __func__, n); - if (owner.count && pss->ac_alloc_task) { - lwsl_notice("%s: orig exit\n", __func__); - sais_event_db_close(vhd, &pdb); - lwsac_free(&ac); - lwsac_free(&failed_ac); - memcpy(&pss->alloc_task, lws_container_of( - owner.head, sai_task_t, list), sizeof(pss->alloc_task)); + if (nn != SQLITE_ROW) + break; - lwsl_notice("%s: platform %s: returning fallback task\n", __func__, platform); + t = sqlite3_column_text(sm, 0); + if (!t) + continue; - return &pss->alloc_task; + /* + * We found errored tasks in the previous event for this + * repo / branch / platform. Let's record them in a temp + * lwsac and condsider if we should use this info to + * prioritize running the corresponding task in the current + * event first + */ + + fti = lwsac_use_zero(&failed_ac, sizeof(*fti) + + strlen((const char *)t) + 1, 256); + if (fti) { + fti->taskname = (const char *)&fti[1]; + memcpy((char *)fti->taskname, t, + strlen((const char *)t) + 1); + lws_dll2_add_tail(&fti->list, &failed_tasks_owner); + } } - } else - lwsl_notice("%s: platform %s: no pending count\n", __func__, platform); + sqlite3_finalize(sm); + } else + lwsl_err("%s: query fail 1\n", __func__); + + sais_event_db_close(vhd, &prev_pdb); + } + + /* + * Let's go through the tasks that failed last time we built this repo / branch, and see + * if we can find the analagous task in the current event. + */ + lws_start_foreach_dll(struct lws_dll2 *, p_fail, failed_tasks_owner.head) { + sai_failed_task_info_t *fti = lws_container_of(p_fail, sai_failed_task_info_t, list); + char esc_taskname[256]; + lws_dll2_owner_t owner; + + lws_sql_purify(esc_taskname, fti->taskname, sizeof(esc_taskname)); + lws_snprintf(pf, sizeof(pf), + " and (state == 0 or state == 9) and " + "(platform == '%s') and (taskname == '%s')", + esc_plat, esc_taskname); + + lwsac_free(&pss->ac_alloc_task); + lws_dll2_owner_clear(&owner); + n = lws_struct_sq3_deserialize(pdb, pf, NULL, + lsm_schema_sq3_map_task, + &owner, &pss->ac_alloc_task, 0, 1); + if (!owner.count) + goto next1; + + lwsl_notice("%s: Prioritizing failed task for %s ('%s')\n", + __func__, platform, fti->taskname); sais_event_db_close(vhd, &pdb); - } + lwsac_free(&ac); + lwsac_free(&failed_ac); + memcpy(&pss->alloc_task, lws_container_of( + owner.head, sai_task_t, list), + sizeof(pss->alloc_task)); + + return &pss->alloc_task; +next1: ; + } lws_end_foreach_dll(p_fail); + + lwsl_notice("%s: no priority\n", __func__); + + /* We have fallen back to doing tasks earliest-first */ + + lws_snprintf(pf, sizeof(pf), + " and (state = 0 or state = 9) and (platform = '%s')", + esc_plat); + + lwsac_free(&pss->ac_alloc_task); + lws_dll2_owner_t owner; + lws_dll2_owner_clear(&owner); + n = lws_struct_sq3_deserialize(pdb, pf, "uid asc ", + lsm_schema_sq3_map_task, + &owner, &pss->ac_alloc_task, 0, 1); + // lwsl_notice("%s: deser returned %d\n", __func__, n); + if (!owner.count || !pss->ac_alloc_task) + goto close_next; + + lwsl_notice("%s: orig exit\n", __func__); + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + lwsac_free(&failed_ac); + memcpy(&pss->alloc_task, lws_container_of( + owner.head, sai_task_t, list), + sizeof(pss->alloc_task)); + + return &pss->alloc_task; + +close_next: + sais_event_db_close(vhd, &pdb); +next: ; } lws_end_foreach_dll(p); bail: @@ -445,8 +449,6 @@ sais_find_or_add_pending_plat(struct vhd *vhd, const char *name) } lws_end_foreach_dll(p); - // lwsl_notice("%s: ->->->->-> adding %s\n", __func__, name); - /* platform name is new, make an entry in the ac */ sp = lwsac_use_zero(&vhd->ac_plats, sizeof(sais_plat_t) + strlen(name) + 1, 512); @@ -628,7 +630,7 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *sp, lwsl_notice("%s: %s: task %s found for %s\n", __func__, platform_name, task_template->uuid, sp->name); - if (sais_create_and_offer_task_step(vhd, task_template->uuid, 3)) + if (sais_create_and_offer_task_step(vhd, task_template->uuid)) return 1; /* yes, we will offer it to him */ @@ -753,7 +755,7 @@ nope: } int -sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char force) +sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) { char event_uuid[33], esc_uuid[129], *p, *start, url[128], mirror_path[256], update[128]; @@ -770,11 +772,10 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for int ret = -1; inflight = sais_is_task_inflight(vhd, NULL, task_uuid, &ul); - - lwsl_notice("%s: caller %d\n", __func__, force); if (inflight /* && ul->started */) { - lwsl_notice("%s: ~~~~~~~ not continuing %s as listed on inflight\n", __func__, task_uuid); + lwsl_notice("%s: ~~~ not continuing %s as listed on inflight\n", + __func__, task_uuid); return 1; } @@ -841,9 +842,11 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for temp_task->git_repo_url = event->repo_fetchurl; /* find builder */ + sp = sais_builder_from_uuid(vhd, temp_task->builder_name); if (!sp) { - lwsl_warn("%s: bailing as can't find builder from %s\n", __func__, temp_task->builder_name); + lwsl_warn("%s: bailing as can't find builder from %s\n", + __func__, temp_task->builder_name); goto bail; } @@ -902,7 +905,8 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid, char for if (!p) { /* no more steps */ sai_uuid_list_t *u; - lwsl_err("%s: +++++++++++++++++++ determined no more steps after build_step %d for task %s, setting SAIES_SUCCESS\n", + lwsl_err("%s: +++ determined no more steps after " + "build_step %d for task %s, setting SAIES_SUCCESS\n", __func__, build_step, temp_task->uuid); sais_set_task_state(vhd, temp_task->uuid, SAIES_SUCCESS, 0, 0); diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index ae51ea5..9657609 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -464,6 +464,138 @@ sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name) } /* + * "reject" packet from the builder is actually a disposition about the + * offered task, it can also indicate ACCEPTED. + */ + +static int +sais_process_rej(struct vhd *vhd, struct pss *pss, + sai_plat_t *sp, sai_rejection_t *rej) +{ + char event_uuid[33], do_remove_uuid = 0, q[128], esc_uuid[129]; + int n, build_step = -1; + sqlite3 *pdb = NULL; + sai_uuid_list_t *ul; + + switch (rej->reason) { + case SAI_TASK_REASON_ACCEPTED: + lwsl_notice("%s: SAI_TASK_REASON_ACCEPTED: %s\n", + __func__, rej->task_uuid); + + /* start build duration only from first step accepted */ + + sai_task_uuid_to_event_uuid(event_uuid, rej->task_uuid); + if (sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) + break; + + lws_sql_purify(esc_uuid, rej->task_uuid, sizeof(esc_uuid)); + lws_snprintf(q, sizeof(q), + "select build_step from tasks where uuid='%s'", + esc_uuid); + + if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &build_step, + NULL) != SQLITE_OK) + build_step = -1; + + /* + * Bump the build step on the accepted task + */ + + build_step++; + lws_snprintf(q, sizeof(q), + "update tasks set build_step=%d " + "where state != 4 and uuid='%s'", + build_step, esc_uuid); + sqlite3_exec(pdb, q, NULL, NULL, NULL); + + if (build_step == 1) { + pss->first_log_timestamp = (uint64_t)lws_now_secs(); + lws_snprintf(q, sizeof(q), + "update tasks set started=%llu where uuid='%s'", + (unsigned long long)pss->first_log_timestamp, esc_uuid); + + if (sqlite3_exec(pdb, q, NULL, NULL, NULL) != SQLITE_OK) + lwsl_notice("%s: unable to set started\n", __func__); + } + + lwsl_notice("%s: exiting, setting build_step %d\n", __func__, build_step); + + sais_event_db_close(vhd, &pdb); + + if (sais_set_task_state(vhd, + rej->task_uuid, + SAIES_BEING_BUILT, + build_step == 1 ? pss->first_log_timestamp : 0, 0)) + break; + + /* leave the uuid listed as inflight until step completed */ + break; + + case SAI_TASK_REASON_DUPE: + lwsl_notice("%s: SAI_TASK_REASON_DUPE: %s\n", + __func__, rej->task_uuid); + break; + + case SAI_TASK_REASON_BUSY: + lwsl_notice("%s: SAI_TASK_REASON_BUSY: Set busy: %s\n", + __func__, rej->task_uuid); + do_remove_uuid = 1; + sais_plat_busy(sp, 1); + break; + + case SAI_TASK_REASON_DESTROYED: + lwsl_notice("%s: SAI_TASK_REASON_DESTROYED: Clear busy: %s\n", + __func__, rej->task_uuid); + do_remove_uuid = 1; + + if (rej->ecode & SAISPRF_EXIT) { + if ((rej->ecode & 0xff) == 0) { + n = SAIES_STEP_SUCCESS; + lwsl_notice("%s: |||| SAIES_STEP_SUCCESS: %s\n", + __func__, rej->task_uuid); + } else { + n = SAIES_FAIL; + lwsl_notice("%s: |||| SAIES_FAIL: %s\n", + __func__, rej->task_uuid); + } + } else + if (rej->ecode & 0x2000) { + n = SAIES_CANCELLED; + lwsl_notice("%s: |||| SAIES_CANCELLED: %s\n", + __func__, rej->task_uuid); + + } else { + n = SAIES_FAIL; + lwsl_notice("%s: |||| SAIES_STEP_FAIL: %s\n", + __func__, rej->task_uuid); + } + + if (sais_set_task_state(vhd, rej->task_uuid, n, 0, + lws_now_secs() - pss->first_log_timestamp)) + return 1; + + sais_plat_busy(sp, 0); + break; + } + + if (do_remove_uuid && + sais_is_task_inflight(vhd, sp, rej->task_uuid, &ul)) { + lwsl_notice("%s: ### Removing %s from inflight\n", + __func__, rej->task_uuid); + sais_inflight_entry_destroy(ul); + // sais_task_clear_build_and_logs(vhd, rej->task_uuid, 1); + } + + if (rej->reason == SAI_TASK_REASON_DESTROYED) + /* uuid will not be found listed as inflight for this */ + sais_create_and_offer_task_step(vhd, rej->task_uuid); + + sais_list_builders(vhd); + + return 0; +} + +/* * Server received a communication from a builder * * buf is lws callback `in` which has LWS_PRE already set aside @@ -474,7 +606,7 @@ 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, unsigned int ss_flags) { - char event_uuid[33], s[128], esc[96], do_remove_uuid; + char event_uuid[33], s[128], esc[96]; sai_resource_requisition_t *rr; sai_resource_wellknown_t *wk; sai_plat_owner_t *bp_owner; @@ -485,7 +617,6 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b lws_wsmsg_info_t info; sai_rejection_t *rej; sai_resource_t *res; - sai_uuid_list_t *ul; lws_dll2_owner_t o; sai_artifact_t *ap; uint8_t xbuf[2048]; @@ -766,119 +897,13 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b break; } - lwsl_notice("%s: builder %s reports task status update, reason: %d, %s, slots %d, mem %d, sto %d\n", + lwsl_notice("%s: builder %s reports task status update, " + "reason: %d, %s, slots %d, mem %d, sto %d\n", __func__, sp->name, rej->reason, rej->task_uuid, sp->avail_slots, sp->avail_mem_kib, sp->avail_sto_kib); - do_remove_uuid = 0; - - switch (rej->reason) { - case SAI_TASK_REASON_ACCEPTED: - lwsl_notice("%s: SAI_TASK_REASON_ACCEPTED: %s\n", __func__, rej->task_uuid); - /* start build duration only from first step accepted */ - { - char event_uuid[33]; - sqlite3 *pdb = NULL; - int build_step = -1; - - sai_task_uuid_to_event_uuid(event_uuid, rej->task_uuid); - if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { - char q[128], esc_uuid[129]; - - lws_sql_purify(esc_uuid, rej->task_uuid, sizeof(esc_uuid)); - lws_snprintf(q, sizeof(q), - "select build_step from tasks where uuid='%s'", - esc_uuid); - - if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &build_step, - NULL) != SQLITE_OK) - build_step = -1; - - /* - * Bump the build step on the accepted task - */ - - build_step++; - lws_snprintf(q, sizeof(q), "update tasks set build_step=%d where state != 4 and uuid='%s'", - build_step, esc_uuid); - sqlite3_exec(pdb, q, NULL, NULL, NULL); - - if (build_step == 1) { - pss->first_log_timestamp = (uint64_t)lws_now_secs(); - lws_snprintf(q, sizeof(q), - "update tasks set started=%llu where uuid='%s'", - (unsigned long long)pss->first_log_timestamp, esc_uuid); - - if (sqlite3_exec(pdb, q, NULL, NULL, NULL) != SQLITE_OK) - lwsl_notice("%s: unable to set started\n", __func__); - } - - lwsl_notice("%s: exiting, setting build_step %d\n", __func__, build_step); - - sais_event_db_close(vhd, &pdb); - - if (sais_set_task_state(vhd, - rej->task_uuid, - SAIES_BEING_BUILT, - !build_step ? pss->first_log_timestamp : 0, 0)) - break; - } - } - /* leave the uuid listed as inflight until step completed */ - break; - - case SAI_TASK_REASON_DUPE: - lwsl_notice("%s: SAI_TASK_REASON_DUPE: %s\n", __func__, rej->task_uuid); - break; - - case SAI_TASK_REASON_BUSY: - lwsl_notice("%s: SAI_TASK_REASON_BUSY: Set busy: %s\n", __func__, rej->task_uuid); - do_remove_uuid = 1; - sais_plat_busy(sp, 1); - break; - - case SAI_TASK_REASON_DESTROYED: - lwsl_notice("%s: SAI_TASK_REASON_DESTROYED: Clear busy: %s\n", __func__, rej->task_uuid); - do_remove_uuid = 1; - - if (rej->ecode & SAISPRF_EXIT) { - if ((rej->ecode & 0xff) == 0) { - n = SAIES_STEP_SUCCESS; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_SUCCESS: %s\n", __func__, rej->task_uuid); - } else { - n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL: %s\n", __func__, rej->task_uuid); - } - } else - if (rej->ecode & 0x2000) { - n = SAIES_CANCELLED; - lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED: %s\n", __func__, rej->task_uuid); - - } else { - n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_FAIL: %s\n", __func__, rej->task_uuid); - } - - if (sais_set_task_state(vhd, rej->task_uuid, n, 0, - lws_now_secs() - pss->first_log_timestamp)) - goto bail; - - sais_plat_busy(sp, 0); - break; - } - - if (do_remove_uuid && - sais_is_task_inflight(vhd, sp, rej->task_uuid, &ul)) { - lwsl_notice("%s: ### Removing %s from inflight\n", __func__, rej->task_uuid); - sais_inflight_entry_destroy(ul); - // sais_task_clear_build_and_logs(vhd, rej->task_uuid, 1); - } - - if (rej->reason == SAI_TASK_REASON_DESTROYED) - /* uuid will not be found listed as inflight for this */ - sais_create_and_offer_task_step(vhd, rej->task_uuid, 10); - - sais_list_builders(vhd); + if (sais_process_rej(vhd, pss, sp, rej)) + goto bail; lwsac_free(&pss->a.ac); break; diff --git a/src/server/s-ws-web.c b/src/server/s-ws-web.c index ae4955c..a4595a8 100644 --- a/src/server/s-ws-web.c +++ b/src/server/s-ws-web.c @@ -360,10 +360,8 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) case SAIS_WS_WEBSRV_RX_EVENTDELETE: ei = (sai_browse_rx_evinfo_t *)a.dest; - if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) { - lwsl_err("%s: SAIS_WS_WEBSRV_RX_EVENTDELETE: unable to validate id %s\n", __func__, ei->event_hash); + if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) goto soft_error; - } lwsl_notice("%s: eventdelete %s\n", __func__, ei->event_hash); @@ -403,20 +401,19 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) * Only broadcast to builders if the state has changed * from 0 viewers to >0, or from >0 viewers to 0. */ - if (old_viewers_present != m->vhd->viewers_are_present) { - lwsl_notice("%s: Viewer presence changed to %d. Broadcasting to builders.\n", - __func__, m->vhd->viewers_are_present); - lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) { - struct pss *pss_builder = lws_container_of(p, struct pss, same); - sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); - - if (vsend) { - vsend->viewers = m->vhd->viewers_are_present; - lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); - lws_callback_on_writable(pss_builder->wsi); - } - } lws_end_foreach_dll(p); - } + if (old_viewers_present == m->vhd->viewers_are_present) + break; + + lws_start_foreach_dll(struct lws_dll2 *, p, m->vhd->builders.head) { + struct pss *pss_builder = lws_container_of(p, struct pss, same); + sai_viewer_state_t *vsend = calloc(1, sizeof(*vsend)); + + if (vsend) { + vsend->viewers = m->vhd->viewers_are_present; + lws_dll2_add_tail(&vsend->list, &pss_builder->viewer_state_owner); + lws_callback_on_writable(pss_builder->wsi); + } + } lws_end_foreach_dll(p); break; } case SAIS_WS_WEBSRV_RX_REBUILD:
Page fetched 0s ago, creation time: 7ms (vhost etag hits: 0%, cache hits: 0%)