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 / web / CMakeLists.txt
Author[]Andy Green <andy@warmcat.com> 2025-10-10 05:13 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-10-10 09:07 UTC
Tree349220b17a9cbafb4b6c0c73089c44408f9d542a   Raw Patch
 
builder: refactor tx path
builder: refactor tx path
diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index 591d02c..9df46ba 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -35,416 +35,156 @@ static const lws_struct_map_t lsm_schema_json_loadreport[] = { LSM_SCHEMA (sai_load_report_t, NULL, lsm_load_report_members, "com.warmcat.sai.loadreport"), }; -static lws_ss_state_return_t -saib_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) -{ - struct sai_plat_server *spm = (struct sai_plat_server *)userobj; - //struct sai_plat *sp = (struct sai_plat *)spm->sai_plat; - - lwsl_info("%s: len %d, flags: %d\n", __func__, (int)len, flags); - lwsl_hexdump_info(buf, len); - - if (saib_ws_json_rx_builder(spm, buf, len)) - return 1; - - return 0; -} +/* + * This is the only path to send things from builder->server. + * + * It will copy the incoming buffer fragment into a buflist in order. So you + * should dump all your fragments for a message in here one after the other + * and the message will go out uninterrupted. Having this as the only tx path + * allows us to guarantee we won't interrupt the fragment sequencing. + * + * The fragment sizing does not have to be related to ss usage sizing, it can + * be larger and it will be used from the buflist according to what SS wants. + */ -unsigned int -saib_get_spm_log_count(struct sai_plat_server *spm) +int +saib_srv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags) { - unsigned int spm_logs = 0; - - lws_start_foreach_dll(struct lws_dll2 *, d, builder.sai_plat_owner.head) { - sai_plat_t *p = lws_container_of(d, sai_plat_t, sai_plat_list); + struct sai_plat_server *spm = (struct sai_plat_server *)lws_ss_to_user_object(h); + unsigned int *pi = (unsigned int *)((const char *)buf - sizeof(int)); - lws_start_foreach_dll(struct lws_dll2 *, d1, p->nspawn_owner.head) { - struct sai_nspawn *ns = lws_container_of(d1, struct sai_nspawn, list); + *pi = ss_flags; + + // lwsl_ss_notice(h, "Queuing builder -> sai-server"); + // lwsl_hexdump_notice(buf, len); - if (ns->spm == spm && ns->chunk_cache.count) - spm_logs += ns->chunk_cache.count; + if (lws_buflist_append_segment(&spm->bl_to_srv, buf - sizeof(int), len + sizeof(int)) < 0) + lwsl_ss_err(h, "failed to append"); /* still ask to drain */ - } lws_end_foreach_dll(d1); - } lws_end_foreach_dll(d); + if (lws_ss_request_tx(h)) + lwsl_ss_err(h, "failed to request tx"); - return spm_logs; + return 0; } - -/* - * We cover requested tx for any instance of a platform that can takes tasks - * from the same server... it means just by coming here, no particular - * platform / sai_plat is implied... - */ - -static lws_ss_state_return_t -saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, - int *flags) +int +saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, + const lws_struct_map_t *map, + size_t map_entries, void *object) { - struct sai_plat_server *spm = (struct sai_plat_server *)userobj; - uint8_t *start = buf, *end = buf + (*len) - 1, *p = start; - struct ws_capture_chunk *chunk; - struct sai_plat *sp = NULL; lws_struct_serialize_t *js; - lws_dll2_t *star, *walk; - lws_ss_state_return_t r; - struct sai_nspawn *ns; + unsigned int ssf = LWSSS_FLAG_SOM; + uint8_t buf[4096]; size_t w = 0; - int n = 0; - - /* - * Are there some logs to dump? - */ - - if (saib_get_spm_log_count(spm)) - goto send_logs; - - /* - * Any build metrics to process? - */ - - if (spm->build_metric_list.count) { - struct lws_dll2 *d = lws_dll2_get_head(&spm->build_metric_list); - sai_build_metric_t *m = - lws_container_of(d, sai_build_metric_t, list); - - lwsl_notice("%s: issuing build metric\n", __func__); - - js = lws_struct_json_serialize_create(lsm_schema_build_metric, - LWS_ARRAY_SIZE(lsm_schema_build_metric), 0, m); - if (!js) - return -1; - n = (int)lws_struct_json_serialize(js, start, - lws_ptr_diff_size_t(end, start), &w); - lws_struct_json_serialize_destroy(&js); - - lwsl_hexdump_notice(start, w); - - n = (int)w; - - lws_dll2_remove(&m->list); - free(m); - - r = lws_ss_request_tx(spm->ss); - if (r) - return r; - goto sendify; + js = lws_struct_json_serialize_create(map, map_entries, 0, object); + if (!js) { + lwsl_warn("%s: failed to serialize\n", __func__); + return -1; } - /* - * Any builder state updates / rejections to process? - */ - - if (spm->rejection_list.count) { - struct lws_dll2 *d = lws_dll2_get_head(&spm->rejection_list); - struct sai_rejection *rej = - lws_container_of(d, struct sai_rejection, list); - - lwsl_notice("%s: issuing %s\n", __func__, - rej->task_uuid[0] ? "task rejection" : "load update"); - - js = lws_struct_json_serialize_create(lsm_schema_json_task_rej, - LWS_ARRAY_SIZE(lsm_schema_json_task_rej), 0, rej); - if (!js) + do { + switch (lws_struct_json_serialize(js, buf, sizeof(buf), &w)) { + case LSJS_RESULT_CONTINUE: + lwsl_notice("%s: LSJS_RESULT_CONTINUE\n", __func__); + break; + case LSJS_RESULT_FINISH: + lwsl_notice("%s: LSJS_RESULT_FINISH\n", __func__); + ssf |= LWSSS_FLAG_EOM; + break; + case LSJS_RESULT_ERROR: + lwsl_warn("%s: serialization failed\n", __func__); return -1; + } - n = (int)lws_struct_json_serialize(js, start, - lws_ptr_diff_size_t(end, start), &w); - lws_struct_json_serialize_destroy(&js); - - lwsl_hexdump_notice(start, w); - - n = (int)w; - - lws_dll2_remove(&rej->list); - free(rej); - - r = lws_ss_request_tx(spm->ss); - if (r) - return r; - goto sendify; - } - - /* - * Any load reports to send? - */ - if (spm->load_report_owner.count) { - struct lws_dll2 *d = lws_dll2_get_head(&spm->load_report_owner); - sai_load_report_t *lr = - lws_container_of(d, sai_load_report_t, list); - - // lwsl_notice("%s: issuing load report for %s\n", __func__, - // lr->builder_name); + lwsl_notice("%s: queueing %d bytes, ss_flags %d\n", __func__, (int)w, ssf); + lwsl_hexdump_notice(buf, w); - js = lws_struct_json_serialize_create(lsm_schema_json_loadreport, - LWS_ARRAY_SIZE(lsm_schema_json_loadreport), 0, lr); - if (!js) + if (saib_srv_queue_tx(h, buf, w, ssf)) return -1; - n = (int)lws_struct_json_serialize(js, start, - lws_ptr_diff_size_t(end, start), &w); - lws_struct_json_serialize_destroy(&js); - - // lwsl_hexdump_notice(start, w); + ssf &= ~((unsigned int)LWSSS_FLAG_SOM); + } while (!(ssf & LWSSS_FLAG_EOM)); - n = (int)w; + lws_struct_json_serialize_destroy(&js); - lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, - lr->active_tasks.head) { - sai_active_task_info_t *ati = lws_container_of(p, - sai_active_task_info_t, list); - lws_dll2_remove(&ati->list); - free(ati); - } lws_end_foreach_dll_safe(p, p1); - - lws_dll2_remove(&lr->list); - free(lr); - - r = lws_ss_request_tx(spm->ss); - if (r) - return r; - goto sendify; - } - - /* - * Any resource requests / relinquishments to process? - */ - - if (spm->resource_req_list.count) { - struct lws_dll2 *d = lws_dll2_get_head(&spm->resource_req_list); - sai_resource_msg_t *resm; - - resm = lws_container_of(d, sai_resource_msg_t, list); - - n = (int)resm->len; - if (*len > resm->len) - *len = resm->len; - memcpy(buf, resm->msg, *len); - - lws_dll2_remove(&resm->list); - free(resm); - - lwsl_notice("%s: forwarding to server %.*s\n", __func__, - (int)(*len), (const char *)buf); - - r = lws_ss_request_tx(spm->ss); - if (r) - return r; - goto sendify; - } - - switch (spm->phase) { - case PHASE_IDLE: - break; - - default: - - // lwsl_notice("%s: ++++++++++++++++ updating with platform status\n", __func__); - - /* - * Update server with platform status - */ - - lws_start_foreach_dll(struct lws_dll2 *, d, - builder.sai_plat_owner.head) { - sai_plat_t *p = lws_container_of(d, sai_plat_t, sai_plat_list); - - lwsl_notice("%s: &&&&&&&&&&&&&&&&&&&& platform %s, windows %d\n", __func__, - p->name, p->windows); + return 0; +} - } lws_end_foreach_dll(d); +static lws_ss_state_return_t +saib_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)userobj; + //struct sai_plat *sp = (struct sai_plat *)spm->sai_plat; - js = lws_struct_json_serialize_create(lsm_schema_map_plat, - LWS_ARRAY_SIZE(lsm_schema_map_plat), 0, - &builder.sai_plat_owner); - if (!js) { - lwsl_err("%s: ++++++++++++++++++ FAILED to serialize plat\n", __func__); - return -1; - } + lwsl_info("%s: len %d, flags: %d\n", __func__, (int)len, flags); + lwsl_hexdump_info(buf, len); - n = (int)lws_struct_json_serialize(js, start, - lws_ptr_diff_size_t(end, start), &w); - lws_struct_json_serialize_destroy(&js); + if (saib_ws_json_rx_builder(spm, buf, len)) + return 1; - sp = (sai_plat_t *)builder.sai_plat_owner.head; - lwsl_hexdump_notice(start, w); + return 0; +} - *len = w; - spm->phase = PHASE_IDLE; - *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; +/* + * We cover requested tx for any instance of a platform that can takes tasks + * from the same server... it means just by coming here, no particular + * platform / sai_plat is implied... + */ - if (saib_get_spm_log_count(spm)) - return lws_ss_request_tx(spm->ss); +static lws_ss_state_return_t +saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, + int *flags) +{ + struct sai_plat_server *spm = (struct sai_plat_server *)userobj; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&spm->bl_to_srv), depi; + char som, som1, eom, final = 1; + size_t fsl, used; - return LWSSSSRET_OK; + if (!spm->bl_to_srv) { + lwsl_notice("%s: nothing to send from builder -> srv\n", __func__); + return LWSSSSRET_TX_DONT_SEND; } - return 1; - -send_logs: + depi = *pi; /* - * Yes somebody has some logs... since we handle all logs on any - * platform doing tasks for the same server, we have to take care not - * to favour draining logs for any busy tasks over letting others - * getting a chance at the mic. If we just scan for guys with logs - * from the start of the list each time, we will never deal with guys - * far from the list head while anybody closer has logs. - * - * For that reason we remember the last dll2 who wrote logs, and start - * looking for the next nspawn with pending logs after him next time. - * (The remembered dll2 is set to NULL when the ns it is inside is - * destroyed). + * We can only issue *len at a time. * - * That requires statefully rotating through... + * Notice we are getting the stored flags from the START of the fragment each time. + * that means we can still see the right flags stored with the fragment, even if we + * have partially used the buflist frag and are partway through it. * - * platform in platform list : nspawn in platform's nspawn list - * = = - * spm->last_logging_platform : spm->last_logging_nspawn - * - * ... filtered for nspawns associated with our spm / server SS link. - * - * The platforms and nspawns are allocated at conf-time statically. + * Ergo, only something to skip if we are at som=1. And also notice that although + * *pi will be right, after the lws_buflist..._use() api, what it points to has been + * destroyed. So we also dereference *pi into depi for use below. */ - star = NULL; - do { - uint32_t tries = builder.sai_plat_owner.count; - - if (!spm->last_logging_nspawn) { - /* start at the start */ - sp = spm->last_logging_platform = lws_container_of( - builder.sai_plat_owner.head, - sai_plat_t, sai_plat_list); - walk = spm->last_logging_nspawn = sp->nspawn_owner.head; - } else { - /* if we can move on, move on */ - sp = spm->last_logging_platform; - walk = spm->last_logging_nspawn->next; - } - - /* if no more nspawns, try moving to next platform */ - while (!walk && tries--) { - if (!sp->sai_plat_list.next) - /* if no more platforms, wrap around to first */ - sp = spm->last_logging_platform = - lws_container_of( - builder.sai_plat_owner.head, - sai_plat_t, sai_plat_list); - else - sp = spm->last_logging_platform = - lws_container_of( - sp->sai_plat_list.next, - sai_plat_t, sai_plat_list); - - /* use the first nspawn in our new platform */ - walk = sp->nspawn_owner.head; - } - - spm->last_logging_nspawn = walk; - - if (walk == star) - return 1; /* nothing to do */ - - if (!star) /* take first usable one as the starting point */ - star = walk; + fsl = lws_buflist_next_segment_len(&spm->bl_to_srv, NULL); - ns = lws_container_of(walk, struct sai_nspawn, list); - if (spm != ns->spm) - continue; - - if (!ns->chunk_cache.count || !ns->chunk_cache.tail) - continue; - - /* - * We're going to process a chunk - */ - - chunk = lws_container_of(ns->chunk_cache.tail, - struct ws_capture_chunk, list); - - lws_dll2_remove(&chunk->list); - - if (ns->task) { - - n = lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), - "{\"schema\":\"com-warmcat-sai-logs\"," - "\"task_uuid\":\"%s\", \"timestamp\": %llu," - "\"channel\": %d, \"len\": %d, ", - ns->task->uuid, (unsigned long long)lws_now_usecs(), - chunk->stdfd, (int)chunk->len); - - if (ns->finished_when_logs_drained && !ns->chunk_cache.count) { - n += lws_snprintf((char *)p + n, lws_ptr_diff_size_t(end, p) - (unsigned int)n, - "\"finished\":%d,", ns->retcode); - sp = ns->sp; - n += lws_snprintf((char *)p + n, lws_ptr_diff_size_t(end, p) - (unsigned int)n, - "\"avail_slots\":%d,\"avail_mem_kib\":%u,\"avail_sto_kib\":%u,", - (int)(sp->job_limit ? sp->job_limit : 6u) - - ((int)sp->nspawn_owner.count - 1), - saib_get_free_ram_kib(), - saib_get_free_disk_kib(builder.home)); - } - - n += lws_snprintf((char *)p + n, lws_ptr_diff_size_t(end, p) - (unsigned int)n, - "\"log\":\""); - - // puts((const char *)&chunk[1]); - // puts((const char *)start); - - n += lws_b64_encode_string((const char *)&chunk[1], - (int)chunk->len, (char *)&start[n], - (int)lws_ptr_diff(end, p) - n - 5); - - p[n++] = '\"'; - p[n++] = '}'; - p[n] = '\0'; - // puts((const char *)start); - } - - ns->chunk_cache_size -= sizeof(*chunk) + chunk->len; - free(chunk); - - /* - * Here we're just looking at the log situation specifically with THIS ns, - * not for any ns that drains to this server - */ - - if (ns->finished_when_logs_drained && !ns->chunk_cache.count) { - /* - * He's in DONE state, and the draining he was waiting - * for has now happened. - * - * Let's move on to UPLOADING_ARTIFACTS if any, this only - * happens after we sent all the related logs. - */ - lwsl_notice("%s: logs cache drained and empty\n", __func__); - ns->finished_when_logs_drained = 0; - if (ns->state != NSSTATE_FAILED) - saib_set_ns_state(ns, NSSTATE_UPLOADING_ARTIFACTS); - } - - break; - - } while (walk != star); + lws_buflist_fragment_use(&spm->bl_to_srv, NULL, 0, &som, &eom); + if (som) { + fsl -= sizeof(int); + lws_buflist_fragment_use(&spm->bl_to_srv, buf, sizeof(int), &som1, &eom); + } + used = (size_t)lws_buflist_fragment_use(&spm->bl_to_srv, (uint8_t *)buf, *len, &som1, &eom); + if (!used) + return LWSSSSRET_TX_DONT_SEND; -sendify: + if (used < fsl || (depi & LWS_WRITE_NO_FIN)) + final = 0; - *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - *len = (unsigned int)n; + *len = used; + *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); - if (spm->phase != PHASE_IDLE || saib_get_spm_log_count(spm)) { - r = lws_ss_request_tx(spm->ss); - if (r) - return r; - } + lwsl_ss_notice(spm->ss, "Sending %d web->srv: ssflags %d", (int)*len, (int)*flags); + lwsl_hexdump_notice(buf, *len); - if (!n) - return 1; + if (spm->bl_to_srv) + return lws_ss_request_tx(spm->ss); - return LWSSSSRET_OK; + return 0; } static int @@ -476,7 +216,6 @@ cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user) { struct sai_plat_server *spm = (struct sai_plat_server *)user; sai_plat_t *sp = lws_container_of(d, sai_plat_t, sai_plat_list); - struct ws_capture_chunk *cc; lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, sp->nspawn_owner.head) { @@ -494,17 +233,6 @@ cleanup_on_ss_disconnect(struct lws_dll2 *d, void *user) if (ns->op && ns->op->lsp) lws_spawn_piped_kill_child_process(ns->op->lsp); - - /* clean up any capture chunks */ - - lws_start_foreach_dll_safe(struct lws_dll2 *, e, e1, - ns->chunk_cache.head) { - cc = lws_container_of(e, - struct ws_capture_chunk, list); - lws_dll2_remove(&cc->list); - free(cc); - } lws_end_foreach_dll_safe(e, e1); - } } lws_end_foreach_dll_safe(d, d1); @@ -517,6 +245,7 @@ saib_sul_load_report_cb(struct lws_sorted_usec_list *sul) struct sai_plat_server *spm = lws_container_of(sul, struct sai_plat_server, sul_load_report); char any_platform_on_this_spm_active = 0; + int n; /* * This builder process may have multiple platforms, each with @@ -535,7 +264,8 @@ saib_sul_load_report_cb(struct lws_sorted_usec_list *sul) lws_start_foreach_dll(struct lws_dll2 *, p, builder.sai_plat_owner.head) { struct sai_plat *sp = lws_container_of(p, sai_plat_t, sai_plat_list); sai_plat_server_ref_t *ref = NULL; - sai_load_report_t *lr; + struct lwsac *ac = NULL; + sai_load_report_t lr; char is_active = 0; /* @@ -575,19 +305,17 @@ saib_sul_load_report_cb(struct lws_sorted_usec_list *sul) /* This platform is active for this spm, or just became idle */ - lr = calloc(1, sizeof(*lr)); - if (!lr) - goto around; + memset(&lr, 0, sizeof(lr)); - lws_strncpy(lr->builder_name, sp->name, sizeof(lr->builder_name)); - lr->core_count = saib_get_cpu_count(); - lr->initial_free_ram_kib = saib_get_total_ram_kib(); - lr->initial_free_disk_kib = saib_get_total_disk_kib(builder.home); - lr->reserved_ram_kib = 0; - lr->reserved_disk_kib = 0; - lr->cpu_percent = (unsigned int)saib_get_system_cpu(&builder); - lr->active_steps = 0; - lws_dll2_owner_clear(&lr->active_tasks); + lws_strncpy(lr.builder_name, sp->name, sizeof(lr.builder_name)); + lr.core_count = saib_get_cpu_count(); + lr.initial_free_ram_kib = saib_get_total_ram_kib(); + lr.initial_free_disk_kib = saib_get_total_disk_kib(builder.home); + lr.reserved_ram_kib = 0; + lr.reserved_disk_kib = 0; + lr.cpu_percent = (unsigned int)saib_get_system_cpu(&builder); + lr.active_steps = 0; + lws_dll2_owner_clear(&lr.active_tasks); if (is_active) { lws_start_foreach_dll(struct lws_dll2 *, d, sp->nspawn_owner.head) { @@ -595,30 +323,34 @@ saib_sul_load_report_cb(struct lws_sorted_usec_list *sul) struct sai_nspawn, list); if (ns->spm == spm && ns->state == NSSTATE_EXECUTING_STEPS && ns->task) { - sai_active_task_info_t *ati = malloc(sizeof(*ati)); + sai_active_task_info_t *ati = lwsac_use_zero(&ac, sizeof(*ati), 512); if (ati) { - memset(ati, 0, sizeof(*ati)); lws_strncpy(ati->task_uuid, ns->task->uuid, sizeof(ati->task_uuid)); lws_strncpy(ati->task_name, ns->task->taskname, sizeof(ati->task_name)); - ati->build_step = ns->current_step; - ati->total_steps = ns->build_step_count; - ati->est_peak_mem_kib = ns->task->est_peak_mem_kib; - ati->est_cpu_load_pct = ns->task->est_cpu_load_pct; - ati->est_disk_kib = ns->task->est_disk_kib; - ati->started = ns->task->started; - lws_dll2_add_tail(&ati->list, &lr->active_tasks); - lr->active_steps++; - - lr->reserved_ram_kib += ns->task->est_peak_mem_kib; - lr->reserved_disk_kib += ns->task->est_disk_kib; + ati->build_step = ns->current_step; + ati->total_steps = ns->build_step_count; + ati->est_peak_mem_kib = ns->task->est_peak_mem_kib; + ati->est_cpu_load_pct = ns->task->est_cpu_load_pct; + ati->est_disk_kib = ns->task->est_disk_kib; + ati->started = ns->task->started; + lws_dll2_add_tail(&ati->list, &lr.active_tasks); + lr.active_steps++; + + lr.reserved_ram_kib += ns->task->est_peak_mem_kib; + lr.reserved_disk_kib += ns->task->est_disk_kib; } } } lws_end_foreach_dll(d); } - lws_dll2_add_tail(&lr->list, &spm->load_report_owner); - if (lws_ss_request_tx(spm->ss)) - lwsl_debug("%s: request tx failed\n", __func__); + n = saib_srv_queue_json_fragments_helper(spm->ss, + lsm_schema_json_loadreport, + LWS_ARRAY_SIZE(lsm_schema_json_loadreport), &lr); + + lwsac_free(&ac); + + if (n) + lwsl_warn("%s: failed to queue fragments\n", __func__); ref->was_active = is_active; @@ -717,12 +449,17 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, case LWSSSCS_CONNECTED: lwsl_ss_user(spm->ss, "CONNECTED"); - spm->phase = PHASE_START_ATTACH; /* Initialize the load report SUL timer for this server connection */ lws_sul_schedule(builder.context, 0, &spm->sul_load_report, saib_sul_load_report_cb, 1); - return lws_ss_request_tx(spm->ss); + if (saib_srv_queue_json_fragments_helper(spm->ss, + lsm_schema_map_plat, + LWS_ARRAY_SIZE(lsm_schema_map_plat), + &builder.sai_plat_owner)) + return -1; + + return 0; case LWSSSCS_DISCONNECTED: /* diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index f4ddd56..b6b3ab0 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -46,31 +46,50 @@ static char csep = '\\'; extern struct lws_vhost *builder_vhost; -struct ws_capture_chunk * +int saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel) { - struct ws_capture_chunk *chunk; + char lj[1600]; + int n = 0; if (!ns || !ns->spm) - return NULL; + return 1; + + if (!ns->task) + return 0; - chunk = malloc(sizeof(*chunk) + len); + n = lws_snprintf(lj, sizeof(lj), + "{\"schema\":\"com-warmcat-sai-logs\"," + "\"task_uuid\":\"%s\", \"timestamp\": %llu," + "\"channel\": %d, \"len\": %d, ", + ns->task->uuid, (unsigned long long)lws_now_usecs(), + channel, (int)len); + + if (ns->retcode_set) { + n += lws_snprintf(lj + n, sizeof(lj) - (unsigned int)n, + "\"finished\":%d,", ns->retcode); + n += lws_snprintf(lj + n, sizeof(lj) - (unsigned int)n, + "\"avail_slots\":%d,\"avail_mem_kib\":%u,\"avail_sto_kib\":%u,", + (int)(ns->sp->job_limit ? ns->sp->job_limit : 6u) - + ((int)ns->sp->nspawn_owner.count - 1), + saib_get_free_ram_kib(), + saib_get_free_disk_kib(builder.home)); + } - if (!chunk) - return NULL; + n += lws_snprintf(lj + n, sizeof(lj) - (unsigned int)n, + "\"log\":\""); - memset(chunk, 0, sizeof(*chunk)); - chunk->us = lws_now_usecs(); - chunk->len = len; - chunk->stdfd = (uint8_t)channel; - if (len) - memcpy(&chunk[1], buf, len); + // puts((const char *)&chunk[1]); + // puts((const char *)start); - ns->chunk_cache_size += sizeof(*chunk) + len; + n += lws_b64_encode_string(buf, (int)len, (char *)&lj[n], + (int)sizeof(lj) - n - 5); - lws_dll2_add_head(&chunk->list, &ns->chunk_cache); + lj[n++] = '\"'; + lj[n++] = '}'; + lj[n] = '\0'; - return chunk; + return saib_srv_queue_tx(ns->spm->ss, lj, (size_t)n, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); } static int @@ -137,7 +156,7 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, return -1; } - if (!saib_log_chunk_create(op->ns, buf, len, lws_spawn_get_stdfd(wsi))) + if (saib_log_chunk_create(op->ns, buf, len, lws_spawn_get_stdfd(wsi))) return -1; return lws_ss_request_tx(op->ns->spm->ss) ? -1 : 0; @@ -177,6 +196,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, lwsl_notice("%s: Process TIMED OUT by Sai\n", __func__); exit_code = -1; ns->retcode = SAISPRF_TIMEDOUT; + ns->retcode_set = 1; goto fail; } @@ -184,6 +204,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, lwsl_notice("%s: Process killed by Sai due to spew\n", __func__); exit_code = -1; ns->retcode = SAISPRF_TERMINATED; + ns->retcode_set = 1; goto fail; } @@ -192,6 +213,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, lwsl_notice("%s: Process Exited with exit code %d\n", __func__, si->si_status); exit_code = si->si_status; + ns->retcode_set = 1; ns->retcode = SAISPRF_EXIT | si->si_status; if (ns->user_cancel) ns->retcode = SAISPRF_TERMINATED; @@ -201,6 +223,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, lwsl_notice("%s: Process Terminated by signal %d / %d\n", __func__, si->si_status, si->si_signo); ns->retcode = SAISPRF_SIGNALLED | si->si_signo; + ns->retcode_set = 1; break; default: lwsl_notice("%s: SI code %d\n", __func__, si->si_code); @@ -209,6 +232,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, #else exit_code = si->retcode & 0xff; ns->retcode = SAISPRF_EXIT | exit_code; + ns->retcode_set = 1; #endif if (exit_code) @@ -267,53 +291,49 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, } if (op->spawn) { - sai_build_metric_t *m; + sai_build_metric_t m; + char hash_input[8192]; + unsigned char hash[32]; + struct lws_genhash_ctx ctx; + int n; if (!ns->spm) { lwsl_err("%s: NULL ns->spm", __func__); goto skip; } - m = malloc(sizeof(*m)); - - if (m) { - char hash_input[8192]; - unsigned char hash[32]; - struct lws_genhash_ctx ctx; - int n; - - memset(m, 0, sizeof(*m)); - - lws_snprintf(hash_input, sizeof(hash_input), "%s%s%s%s", - ns->sp->name, op->spawn, - ns->project_name, ns->ref); - - if (lws_genhash_init(&ctx, LWS_GENHASH_TYPE_SHA256) || - lws_genhash_update(&ctx, hash_input, - strlen(hash_input)) || - lws_genhash_destroy(&ctx, hash)) - lwsl_warn("%s: sha256 failed\n", __func__); - else - for (n = 0; n < 32; n++) - lws_snprintf(m->key + (n * 2), 3, - "%02x", hash[n]); - - lws_strncpy(m->builder_name, ns->sp->name, sizeof(m->builder_name)); - lws_strncpy(m->project_name, ns->project_name, sizeof(m->project_name)); - lws_strncpy(m->ref, ns->ref, sizeof(m->ref)); - lws_strncpy(m->task_uuid, ns->task->uuid, sizeof(m->task_uuid)); - m->unixtime = (uint64_t)time(NULL); - m->us_cpu_user = res->us_cpu_user; - m->us_cpu_sys = res->us_cpu_sys; - m->wallclock_us = us_wallclock; - m->peak_mem_rss = peak_mem_bytes; - m->stg_bytes = du.size_in_bytes; - m->parallel = ns->task->parallel; - - lws_dll2_add_tail(&m->list, &ns->spm->build_metric_list); - if (lws_ss_request_tx(ns->spm->ss)) - lwsl_warn("%s: lws_ss_request_tx failed\n", __func__); - } + memset(&m, 0, sizeof(m)); + + lws_snprintf(hash_input, sizeof(hash_input), "%s%s%s%s", + ns->sp->name, op->spawn, + ns->project_name, ns->ref); + + if (lws_genhash_init(&ctx, LWS_GENHASH_TYPE_SHA256) || + lws_genhash_update(&ctx, hash_input, + strlen(hash_input)) || + lws_genhash_destroy(&ctx, hash)) + lwsl_warn("%s: sha256 failed\n", __func__); + else + for (n = 0; n < 32; n++) + lws_snprintf(m.key + (n * 2), 3, + "%02x", hash[n]); + + lws_strncpy(m.builder_name, ns->sp->name, sizeof(m.builder_name)); + lws_strncpy(m.project_name, ns->project_name, sizeof(m.project_name)); + lws_strncpy(m.ref, ns->ref, sizeof(m.ref)); + lws_strncpy(m.task_uuid, ns->task->uuid, sizeof(m.task_uuid)); + m.unixtime = (uint64_t)time(NULL); + m.us_cpu_user = res->us_cpu_user; + m.us_cpu_sys = res->us_cpu_sys; + m.wallclock_us = us_wallclock; + m.peak_mem_rss = peak_mem_bytes; + m.stg_bytes = du.size_in_bytes; + m.parallel = ns->task->parallel; + + if (saib_srv_queue_json_fragments_helper(ns->spm->ss, + lsm_schema_build_metric, + LWS_ARRAY_SIZE(lsm_schema_build_metric), &m)) + return; } skip: @@ -327,12 +347,6 @@ skip: op->spawn = NULL; } -// 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 @@ -340,15 +354,12 @@ skip: saib_log_chunk_create(ns, NULL, 0, 2); - lwsl_notice("%s: ns finished, waiting to drain %d logs\n", - __func__, ns->chunk_cache.count); - - /* - * saib_task_grace(ns) sets ns->finished_when_logs_drained - */ + lwsl_notice("%s: ns finished\n", __func__); saib_task_grace(ns); saib_set_ns_state(ns, NSSTATE_DONE); + if (ns->state != NSSTATE_FAILED) + saib_set_ns_state(ns, NSSTATE_UPLOADING_ARTIFACTS); ns->reap_cb_called = 1; @@ -366,9 +377,6 @@ fail: 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) sets ns->finished_when_logs_drained - */ saib_task_grace(ns); saib_set_ns_state(ns, NSSTATE_FAILED); diff --git a/src/builder/b-power.c b/src/builder/b-power.c index fba3adb..97b130b 100644 --- a/src/builder/b-power.c +++ b/src/builder/b-power.c @@ -293,9 +293,11 @@ sul_idle_cb(lws_sorted_usec_list_t *sul) struct sai_plat_server *spm = lws_container_of(d, struct sai_plat_server, list); - spm->phase = PHASE_START_ATTACH; - if (lws_ss_request_tx(spm->ss)) - lwsl_ss_warn(spm->ss, "Unable to request tx"); + if (saib_srv_queue_json_fragments_helper(spm->ss, + lsm_schema_map_plat, + LWS_ARRAY_SIZE(lsm_schema_map_plat), + &builder.sai_plat_owner)) + return; } lws_end_foreach_dll(d); diff --git a/src/builder/b-private.h b/src/builder/b-private.h index 9cd447b..3d3ab5c 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -214,7 +214,7 @@ saib_task_destroy(struct sai_nspawn *ns); void saib_task_grace(struct sai_nspawn *ns); -struct ws_capture_chunk * +int saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel); int @@ -263,3 +263,11 @@ saib_get_total_disk_kib(const char *path); int saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp, struct lws_vhost **); +int +saib_srv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int ss_flags); + +int +saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, + const lws_struct_map_t *map, + size_t map_entries, void *object); + diff --git a/src/builder/b-refproxy.c b/src/builder/b-refproxy.c index 097a5f0..30f2dae 100644 --- a/src/builder/b-refproxy.c +++ b/src/builder/b-refproxy.c @@ -74,26 +74,17 @@ resproxy_find_by_cookie(struct sai_plat_server *spm, const char *c, size_t clen) static int saib_queue_yield_message(struct sai_plat_server *spm, const char *c, size_t len) { - sai_resource_msg_t *m = - malloc(sizeof(sai_resource_msg_t) + LWS_PRE + 196); + char msg[256]; + size_t jl; - if (!m) - return -1; - - memset(m, 0, sizeof(*m)); - m->msg = ((char *)&m[1]) + LWS_PRE; /* * We just send the cookie to relinquish the leased resources */ - m->len = (size_t)lws_snprintf((char *)m->msg, 196, + jl = (size_t)lws_snprintf(msg, sizeof(msg), "{\"schema\":\"com-warmcat-sai-resource\"," "\"cookie\":\"%.*s\"}", (int)len, c); - lwsl_notice("%s: %s\n", __func__, m->msg); - - lws_dll2_add_tail(&m->list, &spm->resource_req_list); - - return lws_ss_request_tx(spm->ss) ? -1 : 0; + return saib_srv_queue_tx(spm->ss, msg, jl, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); } int @@ -157,7 +148,6 @@ callback_resproxy(struct lws *wsi, enum lws_callback_reasons reason, { struct sai_plat_server *spm = lws_vhost_user(lws_get_vhost(wsi)); struct rppss *pss = (struct rppss *)user; - sai_resource_msg_t *m; const char *p; size_t al; @@ -213,19 +203,9 @@ callback_resproxy(struct lws *wsi, enum lws_callback_reasons reason, return -1; lws_strnncpy(pss->cookie, p, al, sizeof(pss->cookie)); - - m = malloc(sizeof(*m) + LWS_PRE + len); - memset(m, 0, sizeof(*m)); - m->msg = ((const char *)&m[1] + LWS_PRE); - memcpy((char *)m->msg, in, len); - m->len = len; - - lws_dll2_add_tail(&m->list, &spm->resource_req_list); - lws_dll2_add_tail(&pss->list, &spm->resource_pss_list); - /* try to schedule a write */ - return lws_ss_request_tx(spm->ss) ? -1 : 0; + return saib_srv_queue_tx(spm->ss, in, len, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); case LWS_CALLBACK_RAW_WRITEABLE: if (pss->response) { diff --git a/src/builder/b-task.c b/src/builder/b-task.c index 8a7f5f1..7015f08 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -283,7 +283,7 @@ int saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, const char *rej_task_uuid) { - struct sai_rejection *rej; + struct sai_rejection rej; if (!spm) return -1; @@ -291,11 +291,7 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, if (!spm->ss) return 0; - rej = malloc(sizeof(*rej)); - if (!rej) - return -1; - - memset(rej, 0, sizeof(*rej)); + memset(&rej, 0, sizeof(rej)); /* * Queue a builder status update / @@ -306,21 +302,25 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, lwsl_notice("%s: builder %s occupied reject\n", __func__, sp->name); - lws_strncpy(rej->task_uuid, rej_task_uuid, - sizeof(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", + 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); - lws_dll2_add_tail(&rej->list, &spm->rejection_list); + 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; - return lws_ss_request_tx(spm->ss) ? -1 : 0; + return 0; } void @@ -330,7 +330,6 @@ saib_task_destroy(struct sai_nspawn *ns) lwsl_notice("%s: destroying task %s\n", __func__, ns->task ? ns->task->uuid : "null"); - ns->finished_when_logs_drained = 0; lws_sul_cancel(&ns->sul_cleaner); lws_sul_cancel(&ns->sul_task_cancel); @@ -340,8 +339,11 @@ saib_task_destroy(struct sai_nspawn *ns) */ if (ns->spm) { - ns->spm->phase = PHASE_START_ATTACH; - if (lws_ss_request_tx(ns->spm->ss)) + + if (saib_srv_queue_json_fragments_helper(ns->spm->ss, + lsm_schema_map_plat, + LWS_ARRAY_SIZE(lsm_schema_map_plat), + &builder.sai_plat_owner)) return; /* @@ -465,7 +467,6 @@ void saib_task_grace(struct sai_nspawn *ns) { lwsl_err("%s: +++++ starting task %s grace wait\n", __func__, ns->task ? ns->task->uuid : "null"); - 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); } diff --git a/src/common/include/private.h b/src/common/include/private.h index f2b6e9d..aead6b1 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -167,8 +167,6 @@ struct saib_resproxy { }; struct sai_nspawn { - lws_dll2_owner_t chunk_cache; - char inp[512]; char inp_vn[16]; char path[384]; @@ -204,8 +202,6 @@ struct sai_nspawn { uint64_t worst_mem; uint64_t worst_stg; - size_t chunk_cache_size; - const char *server_name; /* sai-server name who triggered this, eg, 'warmcat' */ const char *project_name; /* name of the git project, eg, 'libwebsockets' */ const char *ref; /* remote refname, eg 'server' */ @@ -223,7 +219,7 @@ struct sai_nspawn { uint8_t stdcount; uint8_t term_budget; - uint8_t finished_when_logs_drained:1; + uint8_t retcode_set:1; uint8_t state_changed:1; uint8_t user_cancel:1; uint8_t reap_cb_called:1; @@ -388,16 +384,14 @@ struct sai_plat; typedef struct sai_plat_server { lws_dll2_t list; - lws_dll2_owner_t rejection_list; - lws_dll2_owner_t build_metric_list; - lws_dll2_owner_t resource_req_list; /* sai_resource_msg_t */ lws_dll2_owner_t resource_pss_list; /* so we can find the cookie */ + struct lws_buflist *bl_to_srv; + char resproxy_path[128]; /* for load reporting */ lws_sorted_usec_list_t sul_load_report; - lws_dll2_owner_t load_report_owner; /* sai_load_report_t */ unsigned int viewer_count; const char *url; @@ -409,7 +403,6 @@ typedef struct sai_plat_server { lws_dll2_t *last_logging_nspawn; struct sai_plat *last_logging_platform; - int phase; int refcount; int index; /* used to create unique build dir path */ @@ -614,6 +607,7 @@ extern const lws_struct_map_t lsm_log[10], lsm_artifact[8], lsm_plat_list[1], + lsm_schema_map_plat[1], lsm_task_rej[5], lsm_task_cancel[1], lsm_schema_json_map_can[1], @@ -625,7 +619,8 @@ extern const lws_struct_map_t lsm_schema_rebuild[1], lsm_schema_build_metric[1], lsm_schema_sq3_map_build_metric[1], - lsm_load_report_members[9] + lsm_load_report_members[9], + lsm_schema_json_task_rej[5] ; extern const lws_struct_map_t lsm_build_metric[12]; extern const lws_struct_map_t lsm_plat[10]; diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 81cecfd..92097a4 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -909,6 +909,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)); pss->wsi = wsi; if (sais_ws_json_rx_builder(vhd, pss, in, len)) return -1; diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c index d6abdc5..87fd60a 100644 --- a/src/web/w-websrv.c +++ b/src/web/w-websrv.c @@ -115,7 +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_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. */
Page fetched 0s ago, creation time: 8ms (vhost etag hits: 0%, cache hits: 0%)