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.png
Author[]Andy Green <andy@warmcat.com> 2025-08-27 03:53 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-08-28 05:50 UTC
Treefdde6f1a647e512f0df2ff4a5552ee438a847e34   Raw Patch
 
feat: Implement server-side scheduler for step-by-step task execution
feat: Implement server-side scheduler for step-by-step task execution

Co-developed-by: Gemini 2.5 Pro
diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index 6bd314d..a53ce35 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -285,19 +285,9 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, } ns->current_step++; - if (ns->current_step < ns->build_step_count) { - /* there are more steps, spawn the next one */ - if (ns) - ns->op = NULL; - if (op->spawn) - free(op->spawn); - free(op); - saib_spawn_step(ns); - return; - } - /* all steps succeeded */ - lwsl_notice("%s: all build steps succeeded\n", __func__); + /* step succeeded, wait for next instruction */ + lwsl_notice("%s: step succeeded\n", __func__); if (op->spawn) free(op->spawn); @@ -323,7 +313,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, return; fail: - n = lws_snprintf(s, sizeof(s), "Build step %d FAILED", ns->current_step + 1); + n = lws_snprintf(s, sizeof(s), "Build step %d FAILED\n", ns->current_step + 1); saib_log_chunk_create(ns, s, (size_t)n, 3); saib_task_grace(ns); @@ -386,7 +376,7 @@ static const char * const runscript_first = "export SAI_LOGPROXY_TTY0=%s\n" "export SAI_LOGPROXY_TTY1=%s\n" "set -e\n" - "cd %s/jobs/$SAI_OVN/$SAI_PROJECT\n" + "cd %s/jobs/$SAI_OVN\n" "rm -rf build\n" "%s < /dev/null\n" "exit $?\n" @@ -410,42 +400,38 @@ static const char * const runscript_next = "export SAI_LOGPROXY_TTY0=%s\n" "export SAI_LOGPROXY_TTY1=%s\n" "set -e\n" - "cd %s/jobs/$SAI_OVN/$SAI_PROJECT\n" + "cd %s/jobs/$SAI_OVN\n" "%s < /dev/null\n" "exit $?\n" ; +static const char * const runscript_build = + "#!/bin/bash -x\n" +#if defined(__APPLE__) + "export PATH=/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/sbin:/usr/sbin\n" +#else + "export PATH=/usr/local/bin:$PATH\n" #endif + "export HOME=%s\n" + "export SAI_OVN=%s\n" + "export SAI_PROJECT=%s\n" + "export SAI_REMOTE_REF=%s\n" + "export SAI_INSTANCE_IDX=%d\n" + "export SAI_PARALLEL=%d\n" + "export SAI_BUILDER_RESOURCE_PROXY=%s\n" + "export SAI_LOGPROXY=%s\n" + "export SAI_LOGPROXY_TTY0=%s\n" + "export SAI_LOGPROXY_TTY1=%s\n" + "set -e\n" + "cd %s/jobs/$SAI_OVN/src\n" + "%s < /dev/null\n" + "exit $?\n" +; -int -saib_spawn_step(struct sai_nspawn *ns); - -int -saib_spawn_build(struct sai_nspawn *ns) -{ - const char *p = ns->task->steps; - int n; - - ns->current_step = 0; - ns->build_step_count = 0; - - lwsl_hexdump_err(ns->task->steps, strlen(ns->task->steps)); - - while ((p = strchr(p, '\n'))) { - ns->build_step_count++; - p++; - } - ns->build_step_count++; - - n = lws_snprintf(ns->pending_mirror_log, sizeof(ns->pending_mirror_log), - "Starting build: %d steps\n", ns->build_step_count); - saib_log_chunk_create(ns, ns->pending_mirror_log, (size_t)n, 3); - - return saib_spawn_step(ns); -} +#endif int -saib_spawn_step(struct sai_nspawn *ns) +saib_spawn_script(struct sai_nspawn *ns) { struct lws_spawn_piped_info info; struct saib_opaque_spawn *op; @@ -481,24 +467,7 @@ saib_spawn_step(struct sai_nspawn *ns) #endif char one_step[4096]; - const char *p_build = ns->task->steps, *q; - int step = 0; - - lwsl_hexdump_notice(ns->task->steps, strlen(ns->task->steps)); - - while (step < ns->current_step && (p_build = strchr(p_build, '\n'))) { - p_build++; - step++; - } - - if (p_build) { - q = strchr(p_build, '\n'); - if (q) - lws_strnncpy(one_step, p_build, q - p_build, sizeof(one_step)); - else - lws_strncpy(one_step, p_build, sizeof(one_step)); - } else - one_step[0] = '\0'; + lws_strncpy(one_step, ns->task->script, sizeof(one_step)); #if defined(WIN32) if (_sopen_s(&fd, args, _O_CREAT | _O_TRUNC | _O_WRONLY, @@ -529,8 +498,17 @@ saib_spawn_step(struct sai_nspawn *ns) ns->slp[0].sockpath, ns->slp[1].sockpath, builder.home, ns->inp, one_step); #else + const char *script_template; + + if (ns->current_step == 0) + script_template = runscript_first; + else if (ns->current_step == 1) + script_template = runscript_next; + else + script_template = runscript_build; + n = lws_snprintf(st, sizeof(st), - ns->current_step ? runscript_next : runscript_first, + script_template, builder.home, ns->fsm.ovname, ns->project_name, ns->ref, ns->instance_idx, 1, diff --git a/src/builder/b-private.h b/src/builder/b-private.h index d4ba8cb..4f1c75b 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -172,10 +172,7 @@ int saib_overlay_unmount(struct sai_nspawn *ns); int -saib_spawn_build(struct sai_nspawn *ns); - -int -saib_spawn_step(struct sai_nspawn *ns); +saib_spawn_script(struct sai_nspawn *ns); int saib_prepare_mount(struct sai_builder *b, struct sai_nspawn *ns); diff --git a/src/builder/b-task.c b/src/builder/b-task.c index 97c066a..51ea66a 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -531,17 +531,12 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) sp->nspawn_owner.head) { struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list); - char found = 0; - - n++; - if (!xns->task) - found = 1; - - if (found) { + if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { ns = xns; break; } - + if (!xns->task && !ns) + ns = xns; } lws_end_foreach_dll_safe(d, d1); if (!ns) { @@ -585,7 +580,19 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) ns->ref = task->git_ref; ns->hash = task->git_hash; ns->git_repo_url = task->git_repo_url; + if (ns->task && ns->task->ac_task_container) + 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->spins = 0; + ns->user_cancel = 0; + ns->us_cpu_user = 0; + ns->us_cpu_sys = 0; + ns->worst_mem = 0; + ns->worst_stg = 0; + } ns->spm = spm; /* bind this task to the spm the req came in on */ { @@ -604,8 +611,8 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) ns->fsm.layers[1] = "env"; #endif - lws_snprintf(ns->fsm.ovname, sizeof(ns->fsm.ovname), "%d-%d.%d", - spm->index, ns->sp->index, ns->instance_idx); + lws_snprintf(ns->fsm.ovname, sizeof(ns->fsm.ovname), "%s", + task->uuid); lwsl_notice("%s: server %s\n", __func__, ns->server_name); lwsl_notice("%s: project %s\n", __func__, ns->project_name); @@ -632,20 +639,6 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) if (mkdir(ns->inp, 0755) && errno != EEXIST) goto ebail; -#if defined(WIN32) - n += lws_snprintf(ns->inp + n, sizeof(ns->inp) - (unsigned int)n, "1%c", - csep); - lws_filename_purify_inplace(ns->inp); - if (mkdir(ns->inp, 0755) && errno != EEXIST) - goto ebail; -#endif - - n += lws_snprintf(ns->inp + n, sizeof(ns->inp) - (unsigned int)n, "%s%c", - ns->project_name, csep); - lws_filename_purify_inplace(ns->inp); - if (mkdir(ns->inp, 0755) && errno != EEXIST) - goto ebail; - /* * Create a pending upload dir to mv artifacts into while * we get on with the next job. @@ -664,8 +657,8 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) ns->user_cancel = 0; ns->spins = 0; - if (saib_spawn_build(ns)) { - lwsl_err("%s: saib_spawn_build failed\n", __func__); + if (saib_spawn_script(ns)) { + lwsl_err("%s: saib_spawn_script failed\n", __func__); goto bail; } diff --git a/src/common/include/private.h b/src/common/include/private.h index 96ce6ed..f821db3 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -50,6 +50,7 @@ typedef enum { SAIES_DELETED = 7, SAIES_NOT_READY_FOR_BUILD = 8, + SAIES_STEP_SUCCESS = 9, } sai_event_state_t; enum { @@ -116,7 +117,7 @@ typedef struct { char uuid[65]; char builder_name[96]; char cpack[128]; - char steps[4096]; + char script[4096]; struct lwsac *ac_task_container; @@ -130,6 +131,7 @@ typedef struct { uint64_t duration; int state; int uid; + int build_step; char told_ongoing; } sai_task_t; @@ -522,7 +524,7 @@ extern const lws_struct_map_t lsm_schema_map_ta[1], lsm_schema_map_plat_simple[1], lsm_event[10], - lsm_task[22], + lsm_task[23], lsm_log[7], lsm_artifact[8], lsm_plat_list[1], diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 861b11f..73dc8ea 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -172,7 +172,8 @@ const lws_struct_map_t lsm_task[] = { LSM_STRING_PTR (sai_task_t, git_ref, "git_ref"), LSM_STRING_PTR (sai_task_t, git_hash, "git_hash"), LSM_STRING_PTR (sai_task_t, git_repo_url, "git_repo_url"), - LSM_CARRAY (sai_task_t, steps, "steps"), + LSM_CARRAY (sai_task_t, script, "script"), + LSM_SIGNED (sai_task_t, build_step, "build_step"), }; const lws_struct_map_t lsm_schema_json_map_task[] = { diff --git a/src/server/s-private.h b/src/server/s-private.h index 374e969..44df052 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -180,13 +180,6 @@ typedef struct sais_sqlite_cache { int refcount; } sais_sqlite_cache_t; -typedef struct sai_ongoing_task { - lws_dll2_t list; - char uuid[65]; - lws_usec_t last_log_timestamp; -} sai_ongoing_task_t; - - typedef struct sais_plat { lws_dll2_t list; const char *plat; @@ -205,7 +198,6 @@ struct vhd { struct lws_dll2_owner sai_powers; struct lws_dll2_owner pending_plats; lws_dll2_owner_t powering_up_list; /* sai_powering_up_plat_t */ - lws_dll2_owner_t ongoing_tasks; /* sai_ongoing_task_t */ struct lwsac *ac_plats; @@ -306,6 +298,9 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, const char *cns_name); int +sais_continue_task(struct vhd *vhd, const char *task_uuid); + +int sais_set_task_state(struct vhd *vhd, const char *builder_name, const char *builder_uuid, const char *task_uuid, int state, uint64_t started, uint64_t duration); diff --git a/src/server/s-task.c b/src/server/s-task.c index f202d0d..10934d9 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -160,14 +160,15 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, */ lws_snprintf(update, sizeof(update), - "update tasks set state=%d%s%s%s%s%s%s%s%s where uuid='%s'", + "update tasks set state=%d%s%s%s%s%s%s%s%s%s where uuid='%s'", state, builder_uuid ? ",builder='": "", builder_uuid ? esc1 : "", builder_uuid ? "'" : "", builder_name ? ",builder_name='" : "", builder_name ? esc : "", builder_name ? "'" : "", - esc3, esc4, esc2); + esc3, esc4, state == SAIES_WAITING ? ",build_step=0" : "", + esc2); if (sqlite3_exec((sqlite3 *)e->pdb, update, NULL, NULL, NULL) != SQLITE_OK) { lwsl_err("%s: %s: %s: fail\n", __func__, update, @@ -183,6 +184,11 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, if (state != task_ostate) { + if (state == SAIES_PASSED_TO_BUILDER && + !vhd->sul_activity.list.owner) + lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, + sais_activity_cb, 1 * LWS_US_PER_SEC); + lwsl_notice("%s: seen task %s %d -> %d\n", __func__, task_uuid, task_ostate, state); @@ -267,19 +273,12 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, sais_event_db_close(vhd, (sqlite3 **)&e->pdb); lwsac_free(&ac); - if (state == SAIES_SUCCESS || state == SAIES_FAIL || state == SAIES_CANCELLED) { - sai_ongoing_task_t *ot = NULL; - lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, vhd->ongoing_tasks.head) { - ot = lws_container_of(p, sai_ongoing_task_t, list); - - if (!strcmp(ot->uuid, task_uuid)) { - lws_dll2_remove(&ot->list); - free(ot); - break; - } - } lws_end_foreach_dll_safe(p, p1); + if (state == SAIES_STEP_SUCCESS) { + sais_continue_task(vhd, task_uuid); + state = SAIES_BEING_BUILT; } + return 0; bail: @@ -802,10 +801,7 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, const char *platform_name) { const sai_task_t *task_template; - char esc1[96], esc2[96]; - lws_dll2_owner_t o; sai_task_t *task = NULL; - int n; if (cb->ongoing >= cb->instances) return 1; @@ -818,100 +814,18 @@ sais_allocate_task(struct vhd *vhd, struct pss *pss, sai_plat_t *cb, if (!task_template) return 1; - task = malloc(sizeof(sai_task_t)); - if (!task) { - lwsac_free(&pss->ac_alloc_task); - return -1; - } - *task = *task_template; - lwsl_notice("%s: %s: task found %s\n", __func__, platform_name, cb->name); /* yes, we will offer it to him */ - if (sais_set_task_state(vhd, cb->name, cb->name, task->uuid, + if (sais_set_task_state(vhd, cb->name, cb->name, task_template->uuid, SAIES_PASSED_TO_BUILDER, lws_now_secs(), 0)) goto bail; /* advance the task state first time we get logs */ pss->mark_started = 1; - /* let's get ahold of his event as well */ - - lws_sql_purify(esc1, task->event_uuid, sizeof(esc1)); - lws_snprintf(esc2, sizeof(esc2), " and uuid='%s'", esc1); - n = lws_struct_sq3_deserialize(vhd->server.pdb, esc2, NULL, - lsm_schema_sq3_map_event, &o, - &pss->a.ac, 0, 1); - if (n < 0 || !o.head) - goto bail; - - task->one_event = lws_container_of(o.head, sai_event_t, list); - - task->server_name = pss->server_name; - task->repo_name = task->one_event->repo_name; - task->git_ref = task->one_event->ref; - task->git_hash = task->one_event->hash; - task->git_repo_url = task->one_event->repo_fetchurl; - task->ac_task_container = pss->a.ac; - - if (!task->git_repo_url || !task->git_repo_url[0]) { - lwsl_err("%s: task %s has no repo_fetchurl, failing\n", - __func__, task->uuid); - sais_set_task_state(vhd, NULL, NULL, task->uuid, SAIES_FAIL, 0, 0); - free(task); - lwsac_free(&pss->ac_alloc_task); - - return 1; /* try again for another task */ - } - - lwsl_notice("%s: windows: %d\n", __func__, cb->windows); - - char url[128], mirror_path[256]; - - lws_strncpy(url, task->one_event->repo_fetchurl, sizeof(url)); - lws_filename_purify_inplace(url); - char *q = url; - while (*q) { - if (*q == '/') - *q = '_'; - if (*q == '.') - *q = '_'; - q++; - } - - lws_snprintf(mirror_path, sizeof(mirror_path), "%s", url); - - if (cb->windows) - lws_snprintf(task->steps, sizeof(task->steps), - "git_helper.bat mirror \"%s\" %s %s %s\n" - "git_helper.bat checkout \"%s\" \"build/%s\" %s\n" - "%s", - task->git_repo_url, task->git_ref, task->git_hash, mirror_path, - mirror_path, task->repo_name, task->git_hash, - task->build); - else - lws_snprintf(task->steps, sizeof(task->steps), - "git_helper.sh mirror \"%s\" %s %s %s\n" - "git_helper.sh checkout \"%s\" \"build/%s\" %s\n" - "%s", - task->git_repo_url, task->git_ref, task->git_hash, mirror_path, - mirror_path, task->repo_name, task->git_hash, - task->build); - - /* - * add to server's estimate of builder's ongoing tasks... - */ - cb->ongoing++; - - lwsl_notice("%s: pre pss->issue_task_owner count %d\n", __func__, pss->issue_task_owner.count); - - lws_dll2_add_tail(&task->pending_assign_list, &pss->issue_task_owner); - lws_callback_on_writable(pss->wsi); - - pss->a.ac = NULL; - - sais_platforms_with_tasks_pending(vhd); + sais_continue_task(vhd, task_template->uuid); /* * We are going to leave here with a live pss->a.ac (pointed into by @@ -935,9 +849,10 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) { struct vhd *vhd = lws_container_of(sul, struct vhd, sul_activity); char *p, *start, *end; - sai_ongoing_task_t *ot; lws_usec_t now; int cat, first = 1; + struct lwsac *ac_events = NULL, *ac_tasks = NULL; + lws_dll2_owner_t o_events, o_tasks; p = start = malloc(8192); if (!p) @@ -950,34 +865,220 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) now = lws_now_usecs(); - lws_start_foreach_dll(struct lws_dll2 *, d, vhd->ongoing_tasks.head) { - ot = lws_container_of(d, sai_ongoing_task_t, list); + if (lws_struct_sq3_deserialize(vhd->server.pdb, + " and state != 3 and state != 4 and state != 5 and state != 7", + NULL, lsm_schema_sq3_map_event, &o_events, &ac_events, 0, 100) >= 0 && + o_events.head) { + lws_start_foreach_dll(struct lws_dll2 *, d, o_events.head) { + sai_event_t *e = lws_container_of(d, sai_event_t, list); + sqlite3 *pdb = NULL; + + if (!sais_event_db_ensure_open(vhd, e->uuid, 0, &pdb)) { + if (lws_struct_sq3_deserialize(pdb, + " and (state = 1 or state = 2)", + NULL, lsm_schema_sq3_map_task, &o_tasks, + &ac_tasks, 0, 100) >= 0 && o_tasks.head) { + lws_start_foreach_dll(struct lws_dll2 *, dt, o_tasks.head) { + sai_task_t *t = lws_container_of(dt, sai_task_t, list); + + if (now - (lws_usec_t)(t->last_updated * LWS_US_PER_SEC) > 10 * LWS_US_PER_SEC) + cat = 1; + else if (now - (lws_usec_t)(t->last_updated * LWS_US_PER_SEC) > 3 * LWS_US_PER_SEC) + cat = 2; + else + cat = 3; + + if (!first) + *p++ = ','; + + p += lws_snprintf(p, lws_ptr_diff_size_t(end, p), + "{\"uuid\":\"%s\",\"cat\":%d}", + t->uuid, cat); + first = 0; + } lws_end_foreach_dll(dt); + } + lwsac_free(&ac_tasks); + sais_event_db_close(vhd, &pdb); + } + } lws_end_foreach_dll(d); + } + lwsac_free(&ac_events); + + *p++ = ']'; + *p++ = '}'; + + if (!first) { + sais_websrv_broadcast(vhd->h_ss_websrv, start, + lws_ptr_diff_size_t(p, start)); + lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, + sais_activity_cb, 1 * LWS_US_PER_SEC); + } + + free(start); +} + +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; + sai_task_t *task = NULL, *task_template; + sai_event_t *event; + sai_plat_t *cb; + struct pss *pss; + sqlite3 *pdb = NULL; + int n, build_step; + struct lwsac *ac = NULL; + + memset(event_uuid, 0, sizeof(event_uuid)); + sai_task_uuid_to_event_uuid(event_uuid, 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)); + lws_snprintf(update, sizeof(update), " and uuid='%s'", esc_uuid); + n = lws_struct_sq3_deserialize(pdb, update, NULL, + lsm_schema_sq3_map_task, &o, &ac, 0, 1); + if (n < 0 || !o.head) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&ac); + return -1; + } + + task_template = lws_container_of(o.head, sai_task_t, list); + task = malloc(sizeof(sai_task_t)); + if (!task) { + lwsac_free(&ac); + sais_event_db_close(vhd, &pdb); + return -1; + } + memset(task, 0, sizeof(*task)); + *task = *task_template; + lwsac_free(&ac); + + build_step = task->build_step; + + /* get the event */ + lws_sql_purify(esc_uuid, event_uuid, sizeof(esc_uuid)); + lws_snprintf(update, sizeof(update), " and uuid='%s'", esc_uuid); + n = lws_struct_sq3_deserialize(vhd->server.pdb, update, NULL, + lsm_schema_sq3_map_event, &o_event, + &task->ac_task_container, 0, 1); + if (n < 0 || !o_event.head) { + sais_event_db_close(vhd, &pdb); + free(task); + return -1; + } + + event = lws_container_of(o_event.head, sai_event_t, list); + task->one_event = event; + task->repo_name = event->repo_name; + task->git_ref = event->ref; + task->git_hash = event->hash; + task->git_repo_url = event->repo_fetchurl; + + /* find builder */ + cb = sais_builder_from_uuid(vhd, task->builder_name, __FILE__, __LINE__); + if (!cb) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&task->ac_task_container); + free(task); + return -1; + } + + lws_strncpy(url, task->one_event->repo_fetchurl, sizeof(url)); + lws_filename_purify_inplace(url); + char *q = url; + while (*q) { + if (*q == '/') *q = '_'; + if (*q == '.') *q = '_'; + q++; + } + lws_snprintf(mirror_path, sizeof(mirror_path), "%s", url); - if (now - ot->last_log_timestamp > 10 * LWS_USEC_PER_SEC) - cat = 1; - else if (now - ot->last_log_timestamp > 3 * LWS_US_PER_SEC) - cat = 2; + switch (build_step) { + case 0: /* git mirror */ + if (cb->windows) + lws_snprintf(task->script, sizeof(task->script), + "git_helper.bat mirror \"%s\" %s %s %s", + task->git_repo_url, task->git_ref, task->git_hash, + mirror_path); else - cat = 3; + lws_snprintf(task->script, sizeof(task->script), + "git_helper.sh mirror \"%s\" %s %s %s", + task->git_repo_url, task->git_ref, task->git_hash, + mirror_path); + break; + case 1: /* git checkout */ + if (cb->windows) + lws_snprintf(task->script, sizeof(task->script), + "git_helper.bat checkout \"%s\" src %s", + mirror_path, task->git_hash); + else + lws_snprintf(task->script, sizeof(task->script), + "git_helper.sh checkout \"%s\" src %s", + mirror_path, task->git_hash); + break; + default: + p = start = task->build; + n = 0; + while (n < build_step - 2 && (p = strchr(p, '\n'))) { + p++; + n++; + } + + if (!p) { /* no more steps */ + 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); + free(task); + return 0; + } - if (!first) - *p++ = ','; + start = p; + p = strchr(start, '\n'); + if (p) + *p = '\0'; - p += lws_snprintf(p, lws_ptr_diff_size_t(end, p), - "{\"uuid\":\"%s\",\"cat\":%d}", - ot->uuid, cat); - first = 0; + lws_strncpy(task->script, start, sizeof(task->script)); + break; + } + /* find builder pss */ + pss = NULL; + lws_start_foreach_dll(struct lws_dll2 *, d, vhd->builders.head) { + struct pss *pss_ = lws_container_of(d, struct pss, same); + if (pss_->wsi == cb->wsi) { + pss = pss_; + break; + } } lws_end_foreach_dll(d); - *p++ = ']'; - *p++ = '}'; + if (!pss) { + sais_event_db_close(vhd, &pdb); + lwsac_free(&task->ac_task_container); + free(task); + return -1; + } - if (lws_ptr_diff(p, start) > 48) - sais_websrv_broadcast(vhd->h_ss_websrv, start, lws_ptr_diff_size_t(p, start)); + task->server_name = pss->server_name; - free(start); + lws_dll2_add_tail(&task->pending_assign_list, &pss->issue_task_owner); + lws_callback_on_writable(pss->wsi); + + lws_sql_purify(esc_uuid, task_uuid, sizeof(esc_uuid)); + lws_snprintf(update, sizeof(update), "update tasks set build_step=%d where uuid='%s'", + build_step + 1, esc_uuid); + sqlite3_exec(pdb, update, NULL, NULL, NULL); - lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, - sais_activity_cb, 1 * LWS_US_PER_SEC); + sais_event_db_close(vhd, &pdb); + + return 0; } diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index ded8f52..f4073a2 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -198,6 +198,7 @@ _sais_websrv_broadcast(struct lws_ss_handle *h, void *arg) if (lws_buflist_append_segment(&m->bltx, a->buf, a->len) < 0) { lwsl_warn("%s: buflist append fail\n", __func__); + lws_ss_start_timeout(h, 1); return; } diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 391e262..dc100f3 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -196,7 +196,6 @@ static void sais_log_to_db(struct vhd *vhd, sai_log_t *log) { sais_logcache_pertask_t *lcpt = NULL; - sai_ongoing_task_t *ot; sai_log_t *hlog; /* @@ -245,38 +244,15 @@ sais_log_to_db(struct vhd *vhd, sai_log_t *log) lws_sul_schedule(vhd->context, 0, &vhd->sul_logcache, sais_dump_logs_to_db, 250 * LWS_US_PER_MS); - /* - * Update ongoing task activity - */ - - ot = NULL; - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->ongoing_tasks.head) { - sai_ongoing_task_t *fot = lws_container_of(p, sai_ongoing_task_t, list); - - if (!strcmp(fot->uuid, log->task_uuid)) { - ot = fot; - break; - } - } lws_end_foreach_dll(p); - - if (!ot) { - ot = malloc(sizeof(*ot)); - if (!ot) - return; - memset(ot, 0, sizeof(*ot)); - lws_strncpy(ot->uuid, log->task_uuid, sizeof(ot->uuid)); - lws_dll2_add_tail(&ot->list, &vhd->ongoing_tasks); - } - - ot->last_log_timestamp = lws_now_usecs(); - if (log->channel == 3 && log->log) { /* control channel */ int step; - if (sscanf(log->log, " Step %d:", &step) == 1) { + if (!memcmp(log->log, " Step ", 5)) { char event_uuid[33]; sqlite3 *pdb = NULL; char q[256], esc_uuid[129]; + step = atoi(&log->log[5]); + sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { @@ -707,7 +683,7 @@ bail: log->timestamp - pss->first_log_timestamp)); if (log->finished & SAISPRF_EXIT) { if ((log->finished & 0xff) == 0) - n = SAIES_SUCCESS; + n = SAIES_STEP_SUCCESS; else n = SAIES_FAIL; } else @@ -1218,7 +1194,7 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, * Don't resend the same status over and over */ - if (memcmp(pss->last_power_report, start, lws_ptr_diff_size_t(p, start) + 1)) { + if (strncmp(pss->last_power_report, (const char *)start, lws_ptr_diff_size_t(p, start) + 1)) { diff = 1; memcpy(pss->last_power_report, start, lws_ptr_diff_size_t(p, start) + 1); }
Page fetched 0s ago, creation time: 9ms (vhost etag hits: 0%, cache hits: 0%)