Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / src / power / p-private.h
Author[]Andy Green <andy@warmcat.com> 2026-10-04 19:08 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-10-04 20:05 UTC
Treeaa74a1b745d7a502d97f57141089c8afcb6ac690   Raw Patch
 
server: task steps can't be offered or acted on twice
server: task steps can't be offered or acted on twice

A task's inflight entry was timestamped when step 1 was offered and
never again, but the prune drops entries "offered and not yet
accepted for 30s" by that timestamp.  So for any task older than 30s,
if the 1s prune ran in the gap between offering the next step and its
ACCEPTED arriving, the entry went, the task (STEP_SUCCESS) looked
pending, and the same step was offered a second time.

If that step was quick and the last one, the builder had finished it
and deleted the job dir before the second offer arrived, so it ran it
again in an empty dir; the server, which had already marked the task
succeeded, then failed it:

  >sais> all 6 steps completed, task succeeded
  >saib> Starting task step 6 ===>
  bash: /home/sai/jobs/.../sai-build-script.sh: No such file or directory

 - restart the inflight clock each time a step is offered

 - builders now say which step an ACCEPTED / DESTROYED etc is about,
   and the server ignores one about a step the task isn't at, or on a
   task whose latest run already ended; it no longer moves build_step,
   the task state or the inflight entry for it.  This also stops a
   pause or rebuild-last-step's killed step from marking the task
   CANCELLED.  Builders that don't say are believed as before.

 - don't offer a step for a task whose latest run is finished or
   paused, and look at the latest run rather than the latest that
   isn't FAIL

 - a build ending in '\n' no longer yields an empty extra step after
   the last one: the line walker now agrees with the step count

 - builder refuses an offer of a step past the task's step count

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
diff --git a/src/builder/b-private.h b/src/builder/b-private.h index d6d8cac..d75d243 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -469,7 +469,7 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, int saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, - const char *rej_task_uuid, unsigned int ecode, + const sai_task_t *task, unsigned int ecode, unsigned int reason); int diff --git a/src/builder/b-task.c b/src/builder/b-task.c index e51ffae..bef6bed 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -809,7 +809,7 @@ saib_set_ns_state(struct sai_nspawn *ns, int state) int saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, - const char *rej_task_uuid, unsigned int ecode, + const sai_task_t *task, unsigned int ecode, unsigned int reason) { struct sai_rejection rej; @@ -826,14 +826,17 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, * Queue a builder task status update */ - if (rej_task_uuid) - lws_strncpy(rej.task_uuid, rej_task_uuid, - sizeof(rej.task_uuid)); + lws_strncpy(rej.task_uuid, task->uuid, sizeof(rej.task_uuid)); lws_snprintf(rej.host_platform, sizeof(rej.host_platform), "%s", sp->name); rej.ecode = ecode; rej.reason = (uint8_t)reason; + /* + * Say which step we mean, so the server can tell a report about the + * step the task is at from one about a step it has moved past + */ + rej.step = (unsigned int)task->build_step + 1; if (saib_srv_queue_json_fragments_helper(spm->ss, lsm_schema_json_task_rej, LWS_ARRAY_SIZE(lsm_schema_json_task_rej), &rej)) { @@ -963,7 +966,7 @@ saib_task_destroy(struct sai_nspawn *ns) /* real work idles us only after the settle time */ builder.last_real_us = lws_now_usecs(); - saib_queue_task_status_update(ns->sp, ns->spm, ns->task->uuid, + saib_queue_task_status_update(ns->sp, ns->spm, ns->task, ecode, SAI_TASK_REASON_DESTROYED); /* @@ -1494,6 +1497,27 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, sp->deserialization_ac = a->ac; /* + * A task has build_step_count steps, numbered from 0. Anything past + * that is not a step at all, and by the time we hear of it we've + * already deleted the job dir after the real last step. + */ + + if (task->build_step_count && + task->build_step >= task->build_step_count) { + lwsl_warn("%s: server offered step %d of %d-step task %s, " + "refusing\n", __func__, task->build_step + 1, + task->build_step_count, task->uuid); + if (saib_queue_task_status_update(sp, spm, task, 0, + SAI_TASK_REASON_DUPE)) { + lwsl_notice("TRAP: saib_queue_task_status_update failed (DUPE)\n"); + return -1; + } + saib_reassess_idle_situation(); + + return 0; + } + + /* * Are we willing to take this task step on? * * We may connect to multiple servers and it's asynchronous @@ -1510,7 +1534,7 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { lwsl_warn("%s: server offered task that's already running. State %d, artifacts %d, op %p\n", __func__, xns->state, xns->count_artifacts, xns->op); - if (saib_queue_task_status_update(sp, spm, task->uuid, 0, + if (saib_queue_task_status_update(sp, spm, task, 0, SAI_TASK_REASON_DUPE)) { lwsl_notice("TRAP: saib_queue_task_status_update failed (DUPE)\n"); return -1; @@ -1543,7 +1567,7 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, goto idle_decline; lwsl_warn("%s: builder rejects offered task\n", __func__); - if (saib_queue_task_status_update(sp, spm, task->uuid, 0, + if (saib_queue_task_status_update(sp, spm, task, 0, SAI_TASK_REASON_BUSY)) { lwsl_notice("TRAP: saib_queue_task_status_update failed (BUSY)\n"); return -1; @@ -1847,7 +1871,7 @@ saib_consider_allocating_task(struct sai_plat_server *spm, lws_struct_args_t *a, builder.ram_reserved_kib += ns->res_ram_kib; builder.disk_reserved_kib += ns->res_disk_kib; - if (saib_queue_task_status_update(sp, spm, task->uuid, 0, SAI_TASK_REASON_ACCEPTED)) + if (saib_queue_task_status_update(sp, spm, task, 0, SAI_TASK_REASON_ACCEPTED)) goto bail; #if defined(__APPLE__) @@ -1861,7 +1885,7 @@ idle_decline: * Unlike BUSY, this doesn't tell the server we can't take real tasks */ lwsl_notice("%s: declining idle task %s\n", __func__, task->uuid); - if (saib_queue_task_status_update(sp, spm, task->uuid, 0, + if (saib_queue_task_status_update(sp, spm, task, 0, SAI_TASK_REASON_IDLE_DECLINED)) return -1; saib_reassess_idle_situation(); diff --git a/src/common/include/private.h b/src/common/include/private.h index da1bdb0..4d3632a 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -425,6 +425,7 @@ typedef struct sai_rejection { unsigned int avail_mem_kib; unsigned int avail_sto_kib; unsigned int ecode; + unsigned int step; /* 1-based; 0 = not given */ unsigned char reason; } sai_rejection_t; @@ -1069,7 +1070,7 @@ extern const lws_struct_map_t lsm_artifact[9], lsm_plat_list[1], lsm_schema_map_plat[1], - lsm_task_rej[4], + lsm_task_rej[5], lsm_task_cancel[3], 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 04dcd26..a790faf 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -248,6 +248,7 @@ const lws_struct_map_t lsm_task_rej[] = { LSM_CARRAY (sai_rejection_t, host_platform, "host_platform"), LSM_CARRAY (sai_rejection_t, task_uuid, "task_uuid"), LSM_JO_UNSIGNED (sai_rejection_t, ecode, "ecode"), + LSM_JO_UNSIGNED (sai_rejection_t, step, "step"), LSM_JO_UNSIGNED (sai_rejection_t, reason, "reason"), }; diff --git a/src/server/s-task.c b/src/server/s-task.c index 1b398a2..699dc2c 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -885,7 +885,16 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) __func__, task_uuid); return 1; } + /* + * We're about to offer its next step, and from here until the + * builder accepts it the prune is allowed to decide the offer + * got lost. The entry was listed when step 1 was offered, so + * unless the clock restarts now, any task more than 30s old + * loses its entry the moment the prune runs in this gap, and + * the next scan for pending work offers the same step again. + */ ul->started = 0; + ul->us_time_listed = lws_now_usecs(); } event_uuid[0] = '\0'; @@ -898,11 +907,11 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) // 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 state != 4 and uuid='%s'", esc_uuid); + lws_snprintf(update, sizeof(update), " and uuid='%s'", esc_uuid); n = lws_struct_sq3_deserialize(pdb, update, "run desc", lsm_schema_sq3_map_task, &o, &ac, 0, 1); if (n < 0 || !o.head) { - lwsl_warn("%s: bailing as nothing with state != 4\n", __func__); + lwsl_warn("%s: bailing as no task %s\n", __func__, task_uuid); sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return -1; @@ -910,7 +919,8 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) task_template = lws_container_of(o.head, sai_task_t, list); - if (task_template->state == SAIES_YIELDED) { + switch (task_template->state) { + case SAIES_YIELDED: /* * The builder stopped this idle task's slice, it has no * more steps to offer until s-idle.c starts a new slice @@ -918,6 +928,26 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) sai_event_db_close(&vhd->sqlite3_cache, &pdb); lwsac_free(&ac); return 0; + + case SAIES_SUCCESS: + case SAIES_FAIL: + case SAIES_CANCELLED: + case SAIES_DELETED: + case SAIES_PAUSED: + /* + * Its latest run isn't going anywhere, whatever build_step + * says. Nothing is in flight for it any more, either. + */ + lwsl_notice("%s: not offering %s in state %d\n", __func__, + task_uuid, task_template->state); + if (inflight) + sais_inflight_entry_destroy(ul); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); + lwsac_free(&ac); + return 1; + + default: + break; } /* @@ -1037,7 +1067,12 @@ sais_create_and_offer_task_step(struct vhd *vhd, const char *task_uuid) n++; } - if (!p) { /* no more steps */ + /* + * A build ending in '\n' leaves us at an empty string after + * its last line rather than NULL; that isn't a step, and + * sais_task_build_step_count() doesn't count it as one either + */ + if (!p || !*p) { /* no more steps */ sai_uuid_list_t *u; lwsl_err("%s: +++ determined no more steps after " diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index ce03903..190c363 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -643,6 +643,51 @@ sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name) } /* + * Builders say which step (1-based) an ACCEPTED or DESTROYED is about. It's + * only news if it's about the step the task is at: for ACCEPTED, the one after + * the last one accepted; for DESTROYED, the last one accepted... and either way + * only while the task's latest run is still going. + * + * Anything else is about a step the task has moved past, eg, a second offer of + * a step that completed in the meantime. Acting on it overruns or rewinds + * build_step, and can fail a task that already succeeded. Builders that don't + * say (step 0) are believed, as before. + * + * *build_step is set to the task's build_step, or -1 if we couldn't read it. + */ + +static int +sais_rej_is_stale(sqlite3 *pdb, const sai_rejection_t *rej, int *build_step) +{ + sqlite3_stmt *sm; + int state = -1; + + *build_step = -1; + + if (sqlite3_prepare_v2(pdb, "select build_step,state from tasks where " + "uuid=? order by run desc limit 1", + -1, &sm, NULL) != SQLITE_OK) + return 0; + + sqlite3_bind_text(sm, 1, rej->task_uuid, -1, SQLITE_TRANSIENT); + if (sqlite3_step(sm) == SQLITE_ROW) { + *build_step = sqlite3_column_int(sm, 0); + state = sqlite3_column_int(sm, 1); + } + sqlite3_finalize(sm); + + if (!rej->step || *build_step < 0) + return 0; + + if (state == SAIES_SUCCESS || state == SAIES_FAIL || + state == SAIES_CANCELLED || state == SAIES_DELETED) + return 1; + + return (int)rej->step != *build_step + + (rej->reason == SAI_TASK_REASON_ACCEPTED); +} + +/* * "reject" packet from the builder is actually a disposition about the * offered task, it can also indicate ACCEPTED. */ @@ -652,7 +697,7 @@ 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[384], esc_uuid[129]; - int n, build_step = -1; + int n, build_step = -1, stale = 0; sqlite3 *pdb = NULL; sai_uuid_list_t *ul; @@ -670,14 +715,25 @@ sais_process_rej(struct vhd *vhd, struct pss *pss, 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' order by run desc limit 1", - esc_uuid); + if (sais_rej_is_stale(pdb, rej, &build_step)) { + /* + * Leave build_step, the state and the inflight entry + * alone, they belong to the step the task is really + * at. Whatever this step does, we'll ignore when it + * ends. + */ + lwsl_warn("%s: %s: stale accept of step %u, at %d\n", + __func__, rej->task_uuid, rej->step, + build_step); + sais_task_logf(vhd, rej->task_uuid, + "builder %s started step %u, which the " + "task is no longer at; ignoring that run", + sp->name, rej->step); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); + break; + } - if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &build_step, - NULL) != SQLITE_OK) - build_step = -1; + lws_sql_purify(esc_uuid, rej->task_uuid, sizeof(esc_uuid)); /* * Bump the build step on the accepted task @@ -746,6 +802,36 @@ sais_process_rej(struct vhd *vhd, struct pss *pss, case SAI_TASK_REASON_DESTROYED: lwsl_notice("%s: SAI_TASK_REASON_DESTROYED: Clear busy: %s\n", __func__, rej->task_uuid); + + sai_task_uuid_to_event_uuid(event_uuid, rej->task_uuid); + if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, event_uuid, 0, &pdb)) { + stale = sais_rej_is_stale(pdb, rej, &build_step); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); + } + + if (stale) { + /* + * It's about a step the task isn't at. Either we + * already ignored it starting, or we rewound the task + * under it (pause, rebuild last step) and stopped it. + * + * The task state is not this step's to change, and + * an inflight entry still waiting for an accept is + * the offer of the step the task is really at. One + * whose step had started is this one, though, and + * it's over. + */ + lwsl_warn("%s: %s: stale end of step %u, at %d\n", + __func__, rej->task_uuid, rej->step, + build_step); + if (sais_is_task_inflight(vhd, NULL, rej->task_uuid, + &ul) && ul->started) + sais_inflight_entry_destroy(ul); + sais_plat_busy(sp, 0); + break; + } + do_remove_uuid = 1; if (rej->ecode & SAISPRF_YIELDED) { @@ -826,7 +912,7 @@ sais_process_rej(struct vhd *vhd, struct pss *pss, // sais_task_clear_build_and_logs(vhd, rej->task_uuid, 1); } - if (rej->reason == SAI_TASK_REASON_DESTROYED) + if (rej->reason == SAI_TASK_REASON_DESTROYED && !stale) /* uuid will not be found listed as inflight for this */ sais_create_and_offer_task_step(vhd, rej->task_uuid);
Page fetched 0s ago, creation time: 6ms (vhost etag hits: 0%, cache hits: 0%)