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 / common / c-sqlite3.c
Author[]Andy Green <andy@warmcat.com> 2025-10-22 12:52 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-10-25 12:14 UTC
Tree571085a98adf59f74ac12e8395398a62465bffa3   Raw Patch
 
sul-per-plat: to handle busy
sul-per-plat: to handle busy
diff --git a/assets/sai.js b/assets/sai.js index 4575bae..d305542 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -505,7 +505,7 @@ function hsanitize(s) return s.toString().replace(/[<>"'&]/g, function(chr) { return '&' + table[chr] + ';'; - }).replace(/\n/g, '\n'); + }).replace(/\r\n/g, '\n').replace(/\n/g, '<br>'); } function renderSpreadsheet(tasks) { diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index 3707bae..5a4e611 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -58,7 +58,8 @@ saib_srv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len, unsigned int s // lwsl_ss_notice(h, "Queuing builder -> sai-server"); // lwsl_hexdump_notice(buf, len); - if (lws_buflist_append_segment(&spm->bl_to_srv, buf - sizeof(int), len + sizeof(int)) < 0) + if (lws_buflist_append_segment(&spm->bl_to_srv, (uint8_t *)buf - sizeof(int), + len + sizeof(int)) < 0) lwsl_ss_err(h, "failed to append"); /* still ask to drain */ if (lws_ss_request_tx(h)) @@ -74,7 +75,7 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, { lws_struct_serialize_t *js; unsigned int ssf = LWSSS_FLAG_SOM; - uint8_t buf[1024]; + uint8_t buf[1024 + LWS_PRE]; size_t w = 0; js = lws_struct_json_serialize_create(map, map_entries, 0, object); @@ -84,7 +85,8 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, } do { - switch (lws_struct_json_serialize(js, buf, sizeof(buf), &w)) { + switch (lws_struct_json_serialize(js, buf + LWS_PRE, + sizeof(buf) - LWS_PRE, &w)) { case LSJS_RESULT_CONTINUE: break; case LSJS_RESULT_FINISH: @@ -95,7 +97,7 @@ saib_srv_queue_json_fragments_helper(struct lws_ss_handle *h, return -1; } - if (saib_srv_queue_tx(h, buf, w, ssf)) + if (saib_srv_queue_tx(h, buf + LWS_PRE, w, ssf)) return -1; ssf &= ~((unsigned int)LWSSS_FLAG_SOM); diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index 781a682..b80d7cb 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -49,7 +49,7 @@ extern struct lws_vhost *builder_vhost; int saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel) { - char lj[1600]; + char lj[2600 + LWS_PRE]; int n = 0; if (!ns || !ns->spm) @@ -58,7 +58,7 @@ saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel) if (!ns->task) return 0; - n = lws_snprintf(lj, sizeof(lj), + n = lws_snprintf(lj + LWS_PRE, sizeof(lj) - LWS_PRE, "{\"schema\":\"com-warmcat-sai-logs\"," "\"task_uuid\":\"%s\", \"timestamp\": %llu," "\"channel\": %d, \"len\": %d, ", @@ -66,9 +66,11 @@ saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel) 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, + n += lws_snprintf(lj + LWS_PRE + n, + sizeof(lj) - LWS_PRE - (unsigned int)n, + "\"finished\":%d,", ns->retcode); + n += lws_snprintf(lj + LWS_PRE + n, + sizeof(lj) - LWS_PRE - (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), @@ -77,20 +79,21 @@ saib_log_chunk_create(struct sai_nspawn *ns, void *buf, size_t len, int channel) ns->retcode_set = 0; } - n += lws_snprintf(lj + n, sizeof(lj) - (unsigned int)n, + n += lws_snprintf(lj + LWS_PRE + n, + sizeof(lj) - LWS_PRE - (unsigned int)n, "\"log\":\""); // puts((const char *)&chunk[1]); // puts((const char *)start); - n += lws_b64_encode_string(buf, (int)len, (char *)&lj[n], - (int)sizeof(lj) - n - 5); + n += lws_b64_encode_string(buf, (int)len, (char *)&lj[LWS_PRE + n], + (int)sizeof(lj) - LWS_PRE - n - 5); - lj[n++] = '\"'; - lj[n++] = '}'; - lj[n] = '\0'; + lj[LWS_PRE + n++] = '\"'; + lj[LWS_PRE + n++] = '}'; + lj[LWS_PRE + n] = '\0'; - return saib_srv_queue_tx(ns->spm->ss, lj, (size_t)n, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); + return saib_srv_queue_tx(ns->spm->ss, lj + LWS_PRE, (size_t)n, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); } static int diff --git a/src/common/include/private.h b/src/common/include/private.h index 8da9761..ca5f60c 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -458,6 +458,8 @@ typedef struct sai_plat { lws_dll2_owner_t servers; /* list of sai_plat_server_ref_t */ + lws_sorted_usec_list_t sul_find_jobs; /* server */ + char peer_ip[48]; const char *name; @@ -467,7 +469,9 @@ typedef struct sai_plat { lws_dll2_owner_t nspawn_owner; struct lwsac *deserialization_ac; + struct lws_context *cx; struct lws *wsi; /* server side only */ + void *vhd; lws_dll2_owner_t env_head; char sai_hash[41]; diff --git a/src/server/s-central.c b/src/server/s-central.c index d8171c5..6b1e48e 100644 --- a/src/server/s-central.c +++ b/src/server/s-central.c @@ -152,25 +152,6 @@ sais_central_cb(lws_sorted_usec_list_t *sul) struct vhd *vhd = lws_container_of(sul, struct vhd, sul_central); /* - * For each builder connected to us, see if it can handle a new task, - * and if so, try to select one matching its supported platforms - */ - - lws_start_foreach_dll(struct lws_dll2 *, p, - vhd->server.builder_owner.head) { - sai_plat_t *cb = lws_container_of(p, sai_plat_t, sai_plat_list); - - if (cb->wsi && lws_wsi_user(cb->wsi)) - /* - * try to bind outstanding task to specific builder - * instance - */ - sais_allocate_task(vhd, (struct pss *)lws_wsi_user(cb->wsi), - cb, cb->platform); - - } lws_end_foreach_dll(p); - - /* * Need to globally check for abandoned tasks periodically */ diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 13c4f5a..30ebbca 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -112,6 +112,7 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, struct pss *pss = (struct pss *)user; sai_http_murl_t mu = SHMUT_NONE; const char *pvo_resources, *num; + lws_wsmsg_info_t info; unsigned int ssf; int n; @@ -395,10 +396,15 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, */ n = lws_snprintf((char *)start, sizeof(buf) - LWS_PRE, "{\"schema\":\"sai-overview\"}"); - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, - (const char *)start, (size_t)n, - SAI_WEBSRV_PB__GENERATED, - 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 (lws_return_http_status(wsi, diff --git a/src/server/s-private.h b/src/server/s-private.h index f2c187e..1dba71f 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -31,6 +31,7 @@ struct sai_plat; /* lws_wsmsg_ array for different sources */ enum { SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, + SAI_WEBSRV_PB__PROXIED_FROM_BUILDER_LR, SAI_WEBSRV_PB__LOGS, SAI_WEBSRV_PB__GENERATED, SAI_WEBSRV_PB__ACTIVITY, @@ -129,6 +130,8 @@ struct pss { sai_notification_t sn; struct lws_dll2 same; /* owner: vhd.builders */ + struct lws_buflist *onward_reassembly; + sqlite3 *pdb_artifact; sqlite3_blob *blob_artifact; @@ -327,7 +330,8 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, uint64_t started, uint64_t duration); int -sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv, const char *str, size_t len, int reassembly_idx, unsigned int ss_flags); +sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv, + lws_wsmsg_info_t *info); int sql3_get_integer_cb(void *user, int cols, char **values, char **name); @@ -421,3 +425,12 @@ sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid); int sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl, unsigned int ss_flags); + +void +sais_plat_find_jobs_cb(lws_sorted_usec_list_t *sul); + +void +sais_plat_busy(sai_plat_t *sp, char set); + +void +sais_websrv_broadcast_buflist(struct lws_ss_handle *hsrv, struct lws_buflist **bl); diff --git a/src/server/s-task.c b/src/server/s-task.c index 133882d..7bdcc65 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -915,6 +915,7 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) struct lwsac *ac_events = NULL, *ac_tasks = NULL; lws_dll2_owner_t o_events, o_tasks; char *p, *start, *end, *ast, s = 1; + lws_wsmsg_info_t info; int cat, first = 1; lws_usec_t now; @@ -966,9 +967,13 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) first = 0; if (lws_ptr_diff_size_t(end, p) < 100) { + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__ACTIVITY; + info.buf = (uint8_t *)start; + info.len = lws_ptr_diff_size_t(p, start); + info.ss_flags = s ? LWSSS_FLAG_SOM : 0; /* we might start it, but it won't be the final frag here since we have JSON closure to do */ - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, start, lws_ptr_diff_size_t(p, start), - SAI_WEBSRV_PB__ACTIVITY, (s ? LWSSS_FLAG_SOM : 0)); + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info); p = start; s = 0; } @@ -985,9 +990,12 @@ sais_activity_cb(lws_sorted_usec_list_t *sul) *p++ = '}'; if (!s) { /* ie, if we sent something, send the closing part of the JSON */ - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, start, - lws_ptr_diff_size_t(p, start), - SAI_WEBSRV_PB__ACTIVITY, LWSSS_FLAG_EOM); + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__ACTIVITY; + info.buf = (uint8_t *)start; + info.len = lws_ptr_diff_size_t(p, start); + info.ss_flags = LWSSS_FLAG_EOM; + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info); lws_sul_schedule(vhd->context, 0, &vhd->sul_activity, sais_activity_cb, 1 * LWS_US_PER_SEC); } @@ -1208,3 +1216,39 @@ bail: return ret; } + + +void +sais_plat_find_jobs_cb(lws_sorted_usec_list_t *sul) +{ + sai_plat_t *sp = lws_container_of(sul, sai_plat_t, sul_find_jobs); + + if (!sp->busy && sp->wsi && lws_wsi_user(sp->wsi)) + /* + * try to bind outstanding task to specific builder + * instance + */ + sais_allocate_task((struct vhd *)sp->vhd, + (struct pss *)lws_wsi_user(sp->wsi), + sp, sp->platform); + + if (!sp->busy) + lws_sul_schedule(sp->cx, 0, &sp->sul_find_jobs, + sais_plat_find_jobs_cb, 1 * LWS_US_PER_SEC); +} + +void +sais_plat_busy(sai_plat_t *sp, char set) +{ + if (set) { + lwsl_notice("%s: %s: SETTING BUSY\n", __func__, sp->name); + lws_sul_cancel(&sp->sul_find_jobs); + sp->busy = 1; + return; + } + + sp->busy = 0; + lwsl_notice("%s: %s: CLEARING BUSY\n", __func__, sp->name); + lws_sul_schedule(sp->cx, 0, &sp->sul_find_jobs, + sais_plat_find_jobs_cb, 1 * LWS_US_PER_SEC); +} diff --git a/src/server/s-webops.c b/src/server/s-webops.c index 05f3ac3..72ae57d 100644 --- a/src/server/s-webops.c +++ b/src/server/s-webops.c @@ -51,35 +51,35 @@ * buflist. */ - -typedef struct { - const uint8_t *buf; - size_t len; - unsigned int ss_flags; - int reassembly_idx; -} sais_websrv_broadcast_t; - static void -_sais_websrv_broadcast(struct lws_ss_handle *h, void *v) +_sais_websrv_broadcast(struct lws_ss_handle *h, void *arg) { - websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); - sais_websrv_broadcast_t *a = (sais_websrv_broadcast_t *)v; - unsigned int *pi = (unsigned int *)((const char *)a->buf - sizeof(int)); + websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); + lws_wsmsg_info_t *info = (lws_wsmsg_info_t *)arg; + unsigned int *pi = (unsigned int *)((const char *)info->buf - sizeof(int)); + + info->head_upstream = &m->bl_srv_to_web; + info->private_heads = m->private_heads; + + lwsl_ss_notice(h, "Queueing %u bytes, ridx %d, ff_flags: %u\n", + (unsigned int)info->len, info->private_source_idx, info->ss_flags); - *pi = a->ss_flags; + *pi = info->ss_flags; /* sai-web might not be taking it.. */ if (lws_buflist_total_len(&m->bl_srv_to_web) > (5u * 1024u * 1024u)) { lwsl_ss_warn(h, "server->web buflist reached 5MB"); + /* close the connection to the client then */ lws_ss_start_timeout(h, 1); + return; } - if (lws_wsmsg_append(&m->bl_srv_to_web, - &m->private_heads[a->reassembly_idx], - a->buf - sizeof(int), - a->len + sizeof(int), a->ss_flags) < 0) + info->buf = info->buf - sizeof(int); + info->len = info->len + sizeof(int); + + if (lws_wsmsg_append(info) < 0) lwsl_ss_err(h, "failed to append"); /* still ask to drain */ if (lws_ss_request_tx(h)) @@ -88,21 +88,54 @@ _sais_websrv_broadcast(struct lws_ss_handle *h, void *v) int sais_websrv_broadcast_REQUIRES_LWS_PRE(struct lws_ss_handle *hsrv, - const char *str, size_t len, - int reassembly_idx, unsigned int ss_flags) + lws_wsmsg_info_t *info) { - sais_websrv_broadcast_t a; + /* calls back for every connected client on server */ + lws_ss_server_foreach_client(hsrv, _sais_websrv_broadcast, info); - a.buf = (const uint8_t *)str; /* LWS_PRE behind valid too */ - a.len = len; - a.ss_flags = ss_flags; - a.reassembly_idx = reassembly_idx; + return 0; +} - lws_ss_server_foreach_client(hsrv, _sais_websrv_broadcast, &a); - return 0; +struct sai_bl_args { + uint8_t *buf; + size_t len; +}; + +static void +_sais_websrv_broadcast_buflist(struct lws_ss_handle *h, void *arg) +{ + websrvss_srv_t *m = (websrvss_srv_t *)lws_ss_to_user_object(h); + struct sai_bl_args *sbba = (struct sai_bl_args *)arg; + + if (lws_buflist_append_segment(&m->bl_srv_to_web, sbba->buf, sbba->len) < 0) + lwsl_notice("%s: failed to store buflist segment\n", __func__); } +/* + * We will copy the buflist bl on to every sai-web client connected to our + * sai-server server, then empty bl. + */ + +void +sais_websrv_broadcast_buflist(struct lws_ss_handle *hsrv, struct lws_buflist **bl) +{ + while (*bl) { + struct sai_bl_args sbba; + + sbba.len = lws_buflist_next_segment_len(bl, &sbba.buf); + + lws_ss_server_foreach_client(hsrv, + _sais_websrv_broadcast_buflist, + (void *)&sbba); + + lws_buflist_use_segment(bl, sbba.len); + } + lws_buflist_destroy_all_segments(bl); +} + + + struct sais_arg { const char *uid; int state; @@ -113,6 +146,7 @@ _sais_taskchange(struct lws_ss_handle *h, void *_arg) { struct sais_arg *arg = (struct sais_arg *)_arg; char tc[LWS_PRE + 128], *start = tc + LWS_PRE; + lws_wsmsg_info_t info; int n; n = lws_snprintf(start, sizeof(tc) - LWS_PRE, @@ -120,9 +154,13 @@ _sais_taskchange(struct lws_ss_handle *h, void *_arg) "\"event_hash\":\"%s\", \"state\":%d}", arg->uid, arg->state); - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, start, (size_t)n, - SAI_WEBSRV_PB__GENERATED, - LWSSS_FLAG_SOM | LWSSS_FLAG_EOM) < 0) { + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)start; + info.len = (size_t)n; + info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, &info) < 0) { lwsl_warn("%s: buflist append failed\n", __func__); return; @@ -145,6 +183,7 @@ _sais_eventchange(struct lws_ss_handle *h, void *_arg) { struct sais_arg *arg = (struct sais_arg *)_arg; char tc[LWS_PRE + 128], *start = tc + LWS_PRE; + lws_wsmsg_info_t info; int n; n = lws_snprintf(start, sizeof(tc) - LWS_PRE, @@ -152,9 +191,13 @@ _sais_eventchange(struct lws_ss_handle *h, void *_arg) "\"event_hash\":\"%s\", \"state\":%d}", arg->uid, arg->state); - if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, start, (size_t)n, - SAI_WEBSRV_PB__GENERATED, - LWSSS_FLAG_SOM | LWSSS_FLAG_EOM) < 0) { + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)start; + info.len = (size_t)n; + info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(h, &info) < 0) { lwsl_warn("%s: buflist append failed\n", __func__); return; } @@ -174,9 +217,9 @@ sais_eventchange(struct lws_ss_handle *hsrv, const char *event_uuid, int state) sai_db_result_t sais_event_reset(struct vhd *vhd, const char *event_uuid) { + struct lwsac *ac = NULL; sqlite3 *pdb = NULL; lws_dll2_owner_t o; - struct lwsac *ac = NULL; char *err = NULL; int ret; @@ -229,6 +272,7 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid) { char qu[128], esc[96], pre[LWS_PRE + 128]; struct lwsac *ac = NULL; + lws_wsmsg_info_t info; sqlite3 *pdb = NULL; lws_dll2_owner_t o; char *err = NULL; @@ -287,9 +331,19 @@ sais_event_delete(struct vhd *vhd, const char *event_uuid) sais_event_db_delete_database(vhd, event_uuid); sais_eventchange(vhd->h_ss_websrv, event_uuid, SAIES_DELETED); - len = (size_t)lws_snprintf(pre + LWS_PRE, sizeof(pre) - LWS_PRE, "{\"schema\":\"sai-overview\"}"); - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, pre + LWS_PRE, len, - SAI_WEBSRV_PB__GENERATED, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); + len = (size_t)lws_snprintf(pre + LWS_PRE, sizeof(pre) - LWS_PRE, + "{\"schema\":\"sai-overview\"}"); + + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)pre + LWS_PRE; + info.len = len; + info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) { + lwsl_err("%s: unable to broadcast\n", __func__); + return SAI_DB_RESULT_ERROR; + } return SAI_DB_RESULT_OK; } diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index d4b2760..d5017f9 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -169,6 +169,7 @@ sais_list_builders(struct vhd *vhd) sai_plat_t *builder_from_db; lws_struct_serialize_t *js; struct lwsac *ac = NULL; + lws_wsmsg_info_t info; size_t w; memset(&db_builders_owner, 0, sizeof(db_builders_owner)); @@ -241,11 +242,16 @@ sais_list_builders(struct vhd *vhd) case LSJS_RESULT_FINISH: /* fallthru */ case LSJS_RESULT_CONTINUE: - sais_websrv_broadcast_REQUIRES_LWS_PRE( - vhd->h_ss_websrv, start, - lws_ptr_diff_size_t(p, start), - SAI_WEBSRV_PB__GENERATED, - ss_flags); + memset(&info, 0, sizeof(info)); + + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)start; + info.len = lws_ptr_diff_size_t(p, start); + info.ss_flags = ss_flags; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: unable to broadcast to web\n", __func__); + p = start; ss_flags &= ~((unsigned int)LWSSS_FLAG_SOM); break; @@ -263,9 +269,15 @@ sais_list_builders(struct vhd *vhd) ss_flags |= LWSSS_FLAG_EOM; p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, start, - lws_ptr_diff_size_t(p, start), - SAI_WEBSRV_PB__GENERATED, ss_flags); + memset(&info, 0, sizeof(info)); + + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = (uint8_t *)start; + info.len = lws_ptr_diff_size_t(p, start); + info.ss_flags = ss_flags; + + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: unable to broadcast to web\n", __func__); // lwsl_notice("%s: Broadcasting builder list: %s\n", __func__, start); lwsac_free(&ac); @@ -510,7 +522,7 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, int *flags) { websrvss_srv_t *m = (websrvss_srv_t *)userobj; - int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->bl_srv_to_web), depi; + int *pi = (int *)lws_buflist_get_frag_start_or_NULL(&m->bl_srv_to_web), depi, fl; char som, som1, eom, final = 1; size_t fsl, used; @@ -549,9 +561,26 @@ websrvss_ws_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, final = 0; *len = used; - *flags = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); + fl = (som ? LWSSS_FLAG_SOM : 0) | (final ? LWSSS_FLAG_EOM : 0); + + // lwsl_ss_notice(m->ss, "Sending %d srv->web: ssflags %d", (int)*len, fl); + + if ((fl & LWSSS_FLAG_SOM) && (((*flags) & 3) == 2)) { + lwsl_ss_err(m->ss, "TX: Illegal LWSSS_FLAG_SOM after previous frame without LWSSS_FLAG_EOM"); + assert(0); + } + if (!(fl & LWSSS_FLAG_SOM) && ((*flags) & 3) == 3) { + lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM after previous frame with LWSSS_FLAG_EOM"); + assert(0); + } + if (!(fl & LWSSS_FLAG_SOM) && !((*flags) & 2)) { + lwsl_ss_err(m->ss, "TX: Missing LWSSS_FLAG_SOM on first frame"); + assert(0); + } + + *flags = fl; + - // lwsl_ss_notice(m->ss, "Sending %d srv->web: som %d, som1 %d, depi %d, ssflags %d", (int)*len, som, som1, depi, (int)*flags); // lwsl_hexdump_notice(buf, *len); if (m->bl_srv_to_web) diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 16c2507..bfc8581 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -85,8 +85,9 @@ static void sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) { struct vhd *vhd = lws_container_of(sul, struct vhd, sul_logcache); - sais_logcache_pertask_t *lcpt; char event_uuid[33], sw[192 + LWS_PRE]; + sais_logcache_pertask_t *lcpt; + lws_wsmsg_info_t info; sqlite3 *pdb = NULL; sai_log_t *hlog; char *err; @@ -146,10 +147,16 @@ sais_dump_logs_to_db(lws_sorted_usec_list_t *sul) n = lws_snprintf(sw + LWS_PRE, sizeof(sw) - LWS_PRE, "{\"schema\":\"sai-tasklogs\"," "\"event_hash\":\"%s\"}", lcpt->uuid); - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, sw + LWS_PRE, - (unsigned int)n, - SAI_WEBSRV_PB__LOGS, - LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); + + memset(&info, 0, sizeof(info)); + + info.private_source_idx = SAI_WEBSRV_PB__LOGS; + info.buf = (uint8_t *)sw + LWS_PRE; + info.len = (unsigned int)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: unable to broadcast to web\n", __func__); /* * Destroy the whole task-specific cache, it will regenerate @@ -425,6 +432,7 @@ sais_builder_disconnected(struct vhd *vhd, struct lws *wsi) } lws_dll2_remove(&cb->sai_plat_list); + lws_sul_cancel(&cb->sul_find_jobs); free(cb); // assert(0); @@ -446,6 +454,8 @@ sai_sql3_get_uint64_cb(void *user, int cols, char **values, char **name) * Server received a communication from a builder * * buf is lws callback `in` which has LWS_PRE already set aside + * + * This could contain multiple pieces, including partials concatenated. */ int @@ -458,6 +468,7 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b sai_plat_owner_t *bp_owner; struct lwsac *ac = NULL; sai_plat_t *build, *cb; + lws_wsmsg_info_t info; sai_rejection_t *rej; sai_resource_t *res; sai_uuid_list_t *ul; @@ -476,752 +487,797 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b goto handle; } - /* - * use the schema name on the incoming JSON to decide what kind of - * structure to instantiate - * - * We may have: - * - * - just received a fragment of the whole JSON - * - * - received the JSON and be handling appeneded blob data - */ - - if (!pss->frag) { - memset(&pss->a, 0, sizeof(pss->a)); - pss->a.map_st[0] = lsm_schema_map_ba; - pss->a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_ba); - pss->a.map_st[1] = lsm_schema_map_ba; - pss->a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_map_ba); - pss->a.ac_block_size = 4096; - - lws_struct_json_init_parse(&pss->ctx, NULL, &pss->a); - } else - pss->frag = 0; - - m = lejp_parse(&pss->ctx, (uint8_t *)buf, (int)bl); - - /* - * returns negative, or unused amount... for us, we either had a - * (negative) error, had LEJP_CONTINUE, or if 0/positive, finished - */ - if (m < 0 && m != LEJP_CONTINUE) { - /* an explicit error */ - lwsl_hexdump_err(buf, bl); - lwsl_err("%s: rx JSON decode failed '%s', %d, %s, %s, %d\n", - __func__, lejp_error_to_string(m), m, - pss->ctx.path, pss->ctx.buf, pss->ctx.npos); - lwsac_free(&pss->a.ac); - return 1; - } - - // lwsl_hexdump_notice(buf, bl); - - if (m == LEJP_CONTINUE) { - if (pss->a.top_schema_index == SAIM_WSSCH_BUILDER_LOADREPORT) - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, (const char *)buf, bl, - SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, ss_flags); - pss->frag = 1; - return 0; - } - - if (!pss->a.dest) { - lwsac_free(&pss->a.ac); - lwsl_err("%s: json decode didn't make an object\n", __func__); - return 1; - } - -handle: - switch (pss->a.top_schema_index) { - case SAIM_WSSCH_BUILDER_PLATS: + while (bl) { /* - * builder is sending us an array of platforms it provides us + * use the schema name on the incoming JSON to decide what kind of + * structure to instantiate + * + * We may have: + * + * - just received a fragment of the whole JSON + * + * - received whole JSON + partial of next + * + * - received whole JSONs + * + * - received the JSON and be handling appeneded blob data */ - bp_owner = (sai_plat_owner_t *)pss->a.dest; - - lws_start_foreach_dll(struct lws_dll2 *, pb, - bp_owner->plat_owner.head) { - build = lws_container_of(pb, sai_plat_t, sai_plat_list); - sai_plat_t *live_cb; - - /* - * Step 1: Update this platform in the persistent database. - */ - char q[1024]; - - lws_snprintf(q, sizeof(q), - "INSERT INTO builders (name, platform, last_seen, peer_ip, sai_hash, lws_hash, windows) " - "VALUES ('%s', '%s', %llu, '%s', '%s', '%s', %d) " - "ON CONFLICT(name) DO UPDATE SET last_seen=excluded.last_seen, " - "peer_ip=excluded.peer_ip, sai_hash=excluded.sai_hash, lws_hash=excluded.lws_hash", - build->name, build->platform, (unsigned long long)lws_now_secs(), - pss->peer_ip, build->sai_hash, build->lws_hash, build->windows); - - if (sai_sqlite3_statement(vhd->server.pdb, q, "upsert builder")) - lwsl_err("%s: Failed to upsert builder %s\n", - __func__, build->name); - - /* - * Step 2: Update the long-lived, malloc'd in-memory list. - */ - //cb = sais_builder_from_uuid(vhd, build->name); - //if (cb) - // sais_builder_disconnected(vhd, cb->wsi); - live_cb = sais_builder_from_uuid(vhd, build->name, __FILE__, __LINE__); - if (live_cb) { - /* Already exists (reconnect), just update dynamic info */ - lwsl_err("%s: found live builder for %s\n", __func__, build->name); - live_cb->wsi = pss->wsi; - lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip)); - lws_strncpy(live_cb->sai_hash, build->sai_hash, - sizeof(live_cb->sai_hash)); - lws_strncpy(live_cb->lws_hash, build->lws_hash, - sizeof(live_cb->lws_hash)); - live_cb->windows = build->windows; - live_cb->online = 1; - live_cb->avail_slots = -1; /* ie, unknown */ - live_cb->avail_mem_kib = (unsigned int)-1; - live_cb->avail_sto_kib = (unsigned int)-1; - live_cb->s_avail_slots = live_cb->avail_slots; - live_cb->s_inflight_count = (int)live_cb->inflight_owner.count; - live_cb->s_last_rej_task_uuid[0] = '\0'; - } else { - /* New builder, create a deep-copied, malloc'd object */ - size_t nlen = strlen(build->name) + 1; - size_t plen = strlen(build->platform) + 1; - - lwsl_err("%s: no live for %s\n", __func__, build->name); - - live_cb = malloc(sizeof(*live_cb) + nlen + plen); - if (live_cb) { - char *p_str = (char *)(live_cb + 1); - - memset(live_cb, 0, sizeof(*live_cb)); - live_cb->name = p_str; - memcpy(p_str, build->name, nlen); - live_cb->platform = p_str + nlen; - memcpy(p_str + nlen, build->platform, plen); - lws_strncpy(live_cb->sai_hash, build->sai_hash, - sizeof(live_cb->sai_hash)); - lws_strncpy(live_cb->lws_hash, build->lws_hash, - sizeof(live_cb->lws_hash)); - live_cb->windows = build->windows; - live_cb->avail_slots = 1; /* default */ - live_cb->avail_mem_kib = (unsigned int)-1; - live_cb->avail_sto_kib = (unsigned int)-1; - live_cb->s_avail_slots = live_cb->avail_slots; - live_cb->wsi = pss->wsi; - live_cb->online = 1; - lws_strncpy(live_cb->peer_ip, pss->peer_ip, sizeof(live_cb->peer_ip)); - lws_dll2_add_tail(&live_cb->sai_plat_list, &vhd->server.builder_owner); - } - } + if (!pss->frag) { + memset(&pss->a, 0, sizeof(pss->a)); + pss->a.map_st[0] = lsm_schema_map_ba; + pss->a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_ba); + pss->a.map_st[1] = lsm_schema_map_ba; + pss->a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_map_ba); + pss->a.ac_block_size = 4096; - const char *dot = strchr(build->name, '.'); - if (dot) { - char host[128]; - lws_strnncpy(host, build->name, dot - build->name, sizeof(host)); - sais_set_builder_power_state(vhd, host, 0, 0); - } - } lws_end_foreach_dll(pb); + lws_struct_json_init_parse(&pss->ctx, NULL, &pss->a); + } else + pss->frag = 0; - /* The lwsac from the parsed message is now completely disposable */ - lwsac_free(&pss->a.ac); + m = lejp_parse(&pss->ctx, (uint8_t *)buf, (int)bl); /* - * Now, iterate through the in-memory list of online builders and - * try to allocate a task for each platform that belongs to the - * builder that just connected. - */ - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - cb = lws_container_of(p, sai_plat_t, sai_plat_list); - if (cb->wsi == pss->wsi) { - /* This platform belongs to the connection that sent the message */ - if (sais_allocate_task(vhd, pss, cb, cb->platform) < 0) - goto bail; - } - } lws_end_foreach_dll(p); -#if 0 - lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { - cb = lws_container_of(p, sai_plat_t, sai_plat_list); - if (cb->wsi == pss->wsi) { - /* This platform belongs to the connection that sent the message */ - if (sais_allocate_task(vhd, pss, cb, cb->platform) < 0) - goto bail; - } - } lws_end_foreach_dll(p); -#endif - - /* - * If we did allocate a task in pss->a.ac, responsibility of - * callback_on_writable handler to empty it + * returns negative, or unused amount... for us, we either had a + * (negative) error, had LEJP_CONTINUE, or if 0/positive, finished */ + if (m < 0 && m != LEJP_CONTINUE) { + /* an explicit error */ + lwsl_hexdump_err(buf, bl); + lwsl_err("%s: rx JSON decode failed '%s', %d, %s, %s, %d\n", + __func__, lejp_error_to_string(m), m, + pss->ctx.path, pss->ctx.buf, pss->ctx.npos); + lwsac_free(&pss->a.ac); + return 1; + } - sais_list_builders(vhd); + // lwsl_hexdump_notice(buf, bl); - break; + if (m == LEJP_CONTINUE) { /* ie, we used all of bl and need more */ + if (pss->a.top_schema_index == SAIM_WSSCH_BUILDER_LOADREPORT) { -bail: - lwsac_free(&pss->a.ac); - return -1; + /* + * We can't directly proxy these pieces, because + * with several builders connected and spamming + * fragmented load reports, when we forward them + * the adjacent fragments will be randomly + * ordered. Even though each builder is sending + * them correctly ordered, when all combined + * together on the srv -> web link, the fragments + * will be disorderd. Eg, b1 first frag, b2 + * first frag, b1 last frag, b2 last frag is + * legal for each builder, but illegal when + * proxied and forwarded in the order they were + * received on a single connection. + * + * Instead we have to collect the pieces per- + * builder and forward them when we have an + * atomic message. + */ - case SAIM_WSSCH_BUILDER_LOGS: - /* - * builder is sending us info about task logs - */ + *((unsigned int *)(buf - sizeof(int))) = ss_flags; + if (lws_buflist_append_segment(&pss->onward_reassembly, + buf - sizeof(int), + bl + sizeof(int)) < 0) + return -1; + } - log = (sai_log_t *)pss->a.dest; - sais_log_to_db(vhd, log); + pss->frag = 1; + return 0; + } - if (pss->mark_started) { - pss->mark_started = 0; - pss->first_log_timestamp = log->timestamp; -// if (sais_set_task_state(vhd, NULL, NULL, log->task_uuid, -// SAIES_BEING_BUILT, 0, 0)) -// goto bail; -// sais_create_and_offer_task_step(vhd, log->task_uuid, 11); + if (!pss->a.dest) { + lwsac_free(&pss->a.ac); + lwsl_err("%s: json decode didn't make an object\n", __func__); + return 1; } + handle: - if (log->finished) { - sai_plat_t *cb; - // sai_uuid_list_t *u; - char builder_name[128], esc_uuid[129], q[128], event_uuid[33]; - sqlite3 *pdb = NULL; + // lwsl_notice("%s: bl: %d, m %d, schema: %d\n", __func__, (int)bl, m, pss->a.top_schema_index); + + switch (pss->a.top_schema_index) { + case SAIM_WSSCH_BUILDER_PLATS: /* - * This step is finished, find the builder and update our - * tracking of its state + * builder is sending us an array of platforms it provides us */ - sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); - if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { - builder_name[0] = '\0'; - lws_sql_purify(esc_uuid, log->task_uuid, sizeof(esc_uuid)); - lws_snprintf(q, sizeof(q), - "select builder_name from tasks where uuid='%s'", - esc_uuid); - if (sqlite3_exec(pdb, q, sql3_get_string_cb, builder_name, - NULL) == SQLITE_OK && builder_name[0]) { - cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); - if (cb) { - // sai_uuid_list_t *sul; - - lwsl_notice("%s: builder %s reports step done, slots %d, mem %d, sto %d\n", - __func__, cb->name, log->avail_slots, log->avail_mem_kib, log->avail_sto_kib); - - cb->avail_slots = log->avail_slots; - cb->avail_mem_kib = log->avail_mem_kib; - cb->avail_sto_kib = log->avail_sto_kib; - cb->last_rej_task_uuid[0] = '\0'; - cb->busy = 0; - -#if 0 - lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, cb->inflight_owner.head) { - sul = lws_container_of(d, sai_uuid_list_t, list); - if (!strcmp(sul->uuid, log->task_uuid)) { - sais_inflight_entry_destroy(sul); - break; - } - } lws_end_foreach_dll_safe(d, d1); -#endif - - cb->s_avail_slots = cb->avail_slots; - cb->s_inflight_count = (int)cb->inflight_owner.count; - lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, - sizeof(cb->s_last_rej_task_uuid)); - sais_list_builders(vhd); - } - } - sais_event_db_close(vhd, &pdb); - } + bp_owner = (sai_plat_owner_t *)pss->a.dest; - /* - * We have reached the end of the logs for this task step - */ + lws_start_foreach_dll(struct lws_dll2 *, pb, + bp_owner->plat_owner.head) { + build = lws_container_of(pb, sai_plat_t, sai_plat_list); + sai_plat_t *live_sp; - sais_dump_logs_to_db(&vhd->sul_logcache); + /* + * Step 1: Update this platform in the persistent database. + */ + char q[1024]; -#if 0 - /* - * Remove us from the inflight list - */ + lws_snprintf(q, sizeof(q), + "INSERT INTO builders (name, platform, last_seen, peer_ip, sai_hash, lws_hash, windows) " + "VALUES ('%s', '%s', %llu, '%s', '%s', '%s', %d) " + "ON CONFLICT(name) DO UPDATE SET last_seen=excluded.last_seen, " + "peer_ip=excluded.peer_ip, sai_hash=excluded.sai_hash, lws_hash=excluded.lws_hash", + build->name, build->platform, (unsigned long long)lws_now_secs(), + pss->peer_ip, build->sai_hash, build->lws_hash, build->windows); - if (sais_is_task_inflight(vhd, cb, log->task_uuid, &u)) - sais_inflight_entry_destroy(u); -#endif - - lwsl_notice("%s: \\\\\\\\\\\\\\\\\\ log->finished says 0x%x, dur %lluus\n", - __func__, log->finished, (unsigned long long)( - log->timestamp - pss->first_log_timestamp)); - if (log->finished & SAISPRF_EXIT) { - if ((log->finished & 0xff) == 0) { - n = SAIES_STEP_SUCCESS; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_SUCCESS: %s\n", __func__, log->task_uuid); - } else { - n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL: %s\n", __func__, log->task_uuid); - } - } else - if (log->finished & 0x2000) { - n = SAIES_CANCELLED; - lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED: %s\n", __func__, log->task_uuid); + if (sai_sqlite3_statement(vhd->server.pdb, q, "upsert builder")) + lwsl_err("%s: Failed to upsert builder %s\n", + __func__, build->name); + /* + * Step 2: Update the long-lived, malloc'd in-memory list. + */ + //cb = sais_builder_from_uuid(vhd, build->name); + //if (cb) + // sais_builder_disconnected(vhd, cb->wsi); + live_sp = sais_builder_from_uuid(vhd, build->name, __FILE__, __LINE__); + if (live_sp) { + /* Already exists (reconnect), just update dynamic info */ + lwsl_err("%s: found live builder for %s\n", __func__, build->name); + live_sp->wsi = pss->wsi; + live_sp->cx = lws_get_context(pss->wsi); + live_sp->vhd = vhd; + lws_strncpy(live_sp->peer_ip, pss->peer_ip, sizeof(live_sp->peer_ip)); + lws_strncpy(live_sp->sai_hash, build->sai_hash, + sizeof(live_sp->sai_hash)); + lws_strncpy(live_sp->lws_hash, build->lws_hash, + sizeof(live_sp->lws_hash)); + live_sp->windows = build->windows; + live_sp->online = 1; + live_sp->avail_slots = -1; /* ie, unknown */ + live_sp->avail_mem_kib = (unsigned int)-1; + live_sp->avail_sto_kib = (unsigned int)-1; + live_sp->s_avail_slots = live_sp->avail_slots; + live_sp->s_inflight_count = (int)live_sp->inflight_owner.count; + live_sp->s_last_rej_task_uuid[0] = '\0'; } else { - n = SAIES_FAIL; - lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_FAIL: %s\n", __func__, log->task_uuid); + /* New builder, create a deep-copied, malloc'd object */ + size_t nlen = strlen(build->name) + 1; + size_t plen = strlen(build->platform) + 1; + + lwsl_err("%s: no live for %s\n", __func__, build->name); + + live_sp = malloc(sizeof(*live_sp) + nlen + plen); + if (live_sp) { + char *p_str = (char *)(live_sp + 1); + + memset(live_sp, 0, sizeof(*live_sp)); + live_sp->name = p_str; + memcpy(p_str, build->name, nlen); + live_sp->platform = p_str + nlen; + memcpy(p_str + nlen, build->platform, plen); + lws_strncpy(live_sp->sai_hash, build->sai_hash, + sizeof(live_sp->sai_hash)); + lws_strncpy(live_sp->lws_hash, build->lws_hash, + sizeof(live_sp->lws_hash)); + live_sp->windows = build->windows; + live_sp->avail_slots = 1; /* default */ + live_sp->avail_mem_kib = (unsigned int)-1; + live_sp->avail_sto_kib = (unsigned int)-1; + live_sp->s_avail_slots = live_sp->avail_slots; + live_sp->wsi = pss->wsi; + live_sp->cx = lws_get_context(pss->wsi); + live_sp->vhd = vhd; + live_sp->online = 1; + lws_strncpy(live_sp->peer_ip, pss->peer_ip, sizeof(live_sp->peer_ip)); + lws_dll2_add_tail(&live_sp->sai_plat_list, &vhd->server.builder_owner); + } } - if (sais_set_task_state(vhd, NULL, NULL, log->task_uuid, n, 0, - log->timestamp - pss->first_log_timestamp)) - goto bail; - } + lws_sul_schedule(live_sp->cx, 0, &live_sp->sul_find_jobs, + sais_plat_find_jobs_cb, 1 * LWS_US_PER_SEC); - lwsac_free(&pss->a.ac); + const char *dot = strchr(build->name, '.'); + if (dot) { + char host[128]; + lws_strnncpy(host, build->name, dot - build->name, sizeof(host)); + sais_set_builder_power_state(vhd, host, 0, 0); + } + } lws_end_foreach_dll(pb); - break; + /* The lwsac from the parsed message is now completely disposable */ + lwsac_free(&pss->a.ac); - case SAIM_WSSCH_BUILDER_TASKREJ: + /* + * Now, iterate through the in-memory list of online builders and + * try to allocate a task for each platform that belongs to the + * builder that just connected. + */ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { + cb = lws_container_of(p, sai_plat_t, sai_plat_list); + if (cb->wsi == pss->wsi) { + /* This platform belongs to the connection that sent the message */ + if (sais_allocate_task(vhd, pss, cb, cb->platform) < 0) + goto bail; + } + } lws_end_foreach_dll(p); + #if 0 + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->server.builder_owner.head) { + cb = lws_container_of(p, sai_plat_t, sai_plat_list); + if (cb->wsi == pss->wsi) { + /* This platform belongs to the connection that sent the message */ + if (sais_allocate_task(vhd, pss, cb, cb->platform) < 0) + goto bail; + } + } lws_end_foreach_dll(p); + #endif - /* - * builder is updating us about a task status - */ + /* + * If we did allocate a task in pss->a.ac, responsibility of + * callback_on_writable handler to empty it + */ - rej = (sai_rejection_t *)pss->a.dest; + sais_list_builders(vhd); - if (!rej->task_uuid[0]) break; - rej->host_platform[sizeof(rej->host_platform) - 1] = '\0'; - cb = sais_builder_from_uuid(vhd, rej->host_platform, __FILE__, __LINE__); - if (!cb) { - lwsl_info("%s: unknown builder %s rejecting\n", - __func__, rej->host_platform); + bail: lwsac_free(&pss->a.ac); - break; - } + return -1; - cb->avail_slots = rej->avail_slots; - cb->avail_mem_kib = rej->avail_mem_kib; - cb->avail_sto_kib = rej->avail_sto_kib; + case SAIM_WSSCH_BUILDER_LOGS: + /* + * builder is sending us info about task logs + */ - lwsl_notice("%s: builder %s reports task status update, reason: %d, %s, slots %d, mem %d, sto %d\n", - __func__, cb->name, rej->reason, rej->task_uuid, - cb->avail_slots, cb->avail_mem_kib, cb->avail_sto_kib); + log = (sai_log_t *)pss->a.dest; + sais_log_to_db(vhd, log); - do_remove_uuid = 0; + if (pss->mark_started) { + pss->mark_started = 0; + pss->first_log_timestamp = log->timestamp; + // if (sais_set_task_state(vhd, NULL, NULL, log->task_uuid, + // SAIES_BEING_BUILT, 0, 0)) + // goto bail; + // sais_create_and_offer_task_step(vhd, log->task_uuid, 11); + } - switch (rej->reason) { - case SAI_TASK_REASON_ACCEPTED: - lwsl_notice("%s: SAI_TASK_REASON_ACCEPTED: %s\n", __func__, rej->task_uuid); - { - char event_uuid[33]; + + if (log->finished) { + sai_plat_t *cb; + // sai_uuid_list_t *u; + char builder_name[128], esc_uuid[129], q[128], 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]; + /* + * This step is finished, find the builder. + * + * We don't move on its state until we receive the + * SAI_TASK_REASON_ from the "_REJ" message from the + * builder. + */ - lws_sql_purify(esc_uuid, rej->task_uuid, sizeof(esc_uuid)); + sai_task_uuid_to_event_uuid(event_uuid, log->task_uuid); + if (!sais_event_db_ensure_open(vhd, event_uuid, 0, &pdb)) { + builder_name[0] = '\0'; + lws_sql_purify(esc_uuid, log->task_uuid, sizeof(esc_uuid)); lws_snprintf(q, sizeof(q), - "select build_step from tasks where uuid='%s'", + "select builder_name from tasks where uuid='%s'", esc_uuid); - if (sqlite3_exec(pdb, q, sql3_get_integer_cb, &build_step, - NULL) != SQLITE_OK) - build_step = -1; + if (sqlite3_exec(pdb, q, sql3_get_string_cb, builder_name, + NULL) == SQLITE_OK && builder_name[0]) { + cb = sais_builder_from_uuid(vhd, builder_name, __FILE__, __LINE__); + if (cb) { + // sai_uuid_list_t *sul; + + lwsl_notice("%s: builder %s reports step done, slots %d, mem %d, sto %d\n", + __func__, cb->name, log->avail_slots, log->avail_mem_kib, log->avail_sto_kib); + + cb->avail_slots = log->avail_slots; + cb->avail_mem_kib = log->avail_mem_kib; + cb->avail_sto_kib = log->avail_sto_kib; + cb->last_rej_task_uuid[0] = '\0'; + + cb->s_avail_slots = cb->avail_slots; + cb->s_inflight_count = (int)cb->inflight_owner.count; + lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, + sizeof(cb->s_last_rej_task_uuid)); + sais_list_builders(vhd); + } + } sais_event_db_close(vhd, &pdb); } - if (build_step == 0) - pss->first_log_timestamp = (uint64_t)lws_now_usecs(); + /* + * We have reached the end of the logs for this task step + */ + + sais_dump_logs_to_db(&vhd->sul_logcache); + + lwsl_notice("%s: \\\\\\\\\\\\\\\\\\ log->finished says 0x%x, dur %lluus\n", + __func__, log->finished, (unsigned long long)( + log->timestamp - pss->first_log_timestamp)); + if (log->finished & SAISPRF_EXIT) { + if ((log->finished & 0xff) == 0) { + n = SAIES_STEP_SUCCESS; + lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_SUCCESS: %s\n", __func__, log->task_uuid); + } else { + n = SAIES_FAIL; + lwsl_notice("%s: |||||||||||||||||||| SAIES_FAIL: %s\n", __func__, log->task_uuid); + } + } else + if (log->finished & 0x2000) { + n = SAIES_CANCELLED; + lwsl_notice("%s: |||||||||||||||||||| SAIES_CANCELLED: %s\n", __func__, log->task_uuid); + + } else { + n = SAIES_FAIL; + lwsl_notice("%s: |||||||||||||||||||| SAIES_STEP_FAIL: %s\n", __func__, log->task_uuid); + } + + if (sais_set_task_state(vhd, NULL, NULL, log->task_uuid, n, 0, + log->timestamp - pss->first_log_timestamp)) + goto bail; } - if (sais_set_task_state(vhd, NULL, NULL, rej->task_uuid, - SAIES_BEING_BUILT, 0, 0)) - break; - /* leave the uuid listed 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; - cb->busy = 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; - cb->busy = 0; + lwsac_free(&pss->a.ac); + break; - } - if (do_remove_uuid && - sais_is_task_inflight(vhd, cb, 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); - } + case SAIM_WSSCH_BUILDER_TASKREJ: - 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); + /* + * builder is updating us about a task status + */ - // sais_task_clear_build_and_logs(vhd, rej->task_uuid, 31); + rej = (sai_rejection_t *)pss->a.dest; - cb->s_avail_slots = cb->avail_slots; - cb->s_inflight_count = (int)cb->inflight_owner.count; - lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, - sizeof(cb->s_last_rej_task_uuid)); + if (!rej->task_uuid[0]) + break; - sais_list_builders(vhd); + rej->host_platform[sizeof(rej->host_platform) - 1] = '\0'; + cb = sais_builder_from_uuid(vhd, rej->host_platform, __FILE__, __LINE__); + if (!cb) { + lwsl_info("%s: unknown builder %s rejecting\n", + __func__, rej->host_platform); + lwsac_free(&pss->a.ac); + break; + } - lwsac_free(&pss->a.ac); - break; + cb->avail_slots = rej->avail_slots; + cb->avail_mem_kib = rej->avail_mem_kib; + cb->avail_sto_kib = rej->avail_sto_kib; - case SAIM_WSSCH_BUILDER_LOADREPORT: -// { -// sai_load_report_t *lr = (sai_load_report_t *)pss->a.dest; + lwsl_notice("%s: builder %s reports task status update, reason: %d, %s, slots %d, mem %d, sto %d\n", + __func__, cb->name, rej->reason, rej->task_uuid, + cb->avail_slots, cb->avail_mem_kib, cb->avail_sto_kib); -// lwsl_notice("%s: @@@@@@@@@@@@@@@@@@ loadreport from %s: ram %uk, disk %uk\n", -// __func__, lr->builder_name, lr->reserved_ram_kib, -// lr->reserved_disk_kib); + do_remove_uuid = 0; -// ssize_t wr = write(2, buf, bl); -// if (wr != (ssize_t)bl) -// lwsl_notice("%s: write failed\n", __func__); -// } -// lwsl_wsi_user(pss->wsi, "SAIM_WSSCH_BUILDER_LOADREPORT broadcasting\n"); - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, (const char *)buf, bl, - SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, ss_flags); - break; + 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; - case SAIM_WSSCH_BUILDER_ARTIFACT: - /* - * Builder wants to send us an artifact. - * - * We get sent a JSON object immediately followed by binary - * data for the artifact. - * - * We place the binary data as a blob in the sql record in the - * artifact table. - */ + 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]; - lwsl_info("%s: SAIM_WSSCH_BUILDER_ARTIFACT: m = %d, bl = %d\n", __func__, m, (int)bl); + 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; + sais_event_db_close(vhd, &pdb); + } - if (!pss->bulk_binary_data) { + if (build_step == 0) + pss->first_log_timestamp = (uint64_t)lws_now_usecs(); + } - lwsl_info("%s: BUILDER_ARTIFACT: blob start, m = %d\n", __func__, m); + if (sais_set_task_state(vhd, NULL, NULL, rej->task_uuid, + SAIES_BEING_BUILT, 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(cb, 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; + sais_plat_busy(cb, 0); + break; + } - ap = (sai_artifact_t *)pss->a.dest; + if (do_remove_uuid && + sais_is_task_inflight(vhd, cb, 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); + } - sai_task_uuid_to_event_uuid(event_uuid, ap->task_uuid); + 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); - /* - * Open the event-specific database object... the - * handle is closed when the stream closes, for whatever - * reason. - */ + cb->s_avail_slots = cb->avail_slots; + cb->s_inflight_count = (int)cb->inflight_owner.count; + lws_strncpy(cb->s_last_rej_task_uuid, cb->last_rej_task_uuid, + sizeof(cb->s_last_rej_task_uuid)); - if (sais_event_db_ensure_open(pss->vhd, event_uuid, 0, - &pss->pdb_artifact)) { - lwsl_err("%s: unable to open event-specific " - "database\n", __func__); + sais_list_builders(vhd); - lwsac_free(&pss->a.ac); - return -1; - } + lwsac_free(&pss->a.ac); + break; + + case SAIM_WSSCH_BUILDER_LOADREPORT: /* - * Retreive the task object + * If we got here, we have any intermediate parts + * already, let's add this final part there first */ - lws_sql_purify(esc, ap->task_uuid, sizeof(esc)); - lws_snprintf(s, sizeof(s)," and uuid == \"%s\"", esc); - n = lws_struct_sq3_deserialize(pss->pdb_artifact, s, - NULL, lsm_schema_sq3_map_task, - &o, &ac, 0, 1); - if (n < 0 || !o.head) { - sais_event_db_close(vhd, &pss->pdb_artifact); - lwsl_notice("%s: no task of that id\n", __func__); - lwsac_free(&pss->a.ac); + *((unsigned int *)(buf - sizeof(int))) = ss_flags; + if (lws_buflist_append_segment(&pss->onward_reassembly, + buf - sizeof(int), + bl + sizeof(int)) < 0) return -1; - } - task = (sai_task_t *)o.head; - n = strcmp(task->art_up_nonce, ap->artifact_up_nonce); + /* + * Then let's forward the whole reassembly buflist on + * to the proxying buflist atomically. + */ - if (n) { - lwsl_err("%s: artifact nonce mismatch\n", - __func__); - goto afail; - } + sais_websrv_broadcast_buflist(vhd->h_ss_websrv, + &pss->onward_reassembly); + + break; + case SAIM_WSSCH_BUILDER_ARTIFACT: /* - * The task the sender is sending us an artifact for - * exists. The sender knows the random upload nonce - * for that task's artifacts. + * Builder wants to send us an artifact. * - * Create a random download nonce unrelated to the - * random upload nonce (so knowing the download one - * won't let you upload anything). + * We get sent a JSON object immediately followed by binary + * data for the artifact. * - * Create the artifact's entry in the event-specific - * database + * We place the binary data as a blob in the sql record in the + * artifact table. */ - sai_uuid16_create(pss->vhd->context, - ap->artifact_down_nonce); + lwsl_info("%s: SAIM_WSSCH_BUILDER_ARTIFACT: m = %d, bl = %d\n", __func__, m, (int)bl); - lws_dll2_owner_clear(&o); - lws_dll2_add_head(&ap->list, &o); + if (!pss->bulk_binary_data) { - /* - * Create the task in event-specific database - */ + lwsl_info("%s: BUILDER_ARTIFACT: blob start, m = %d\n", __func__, m); - if (lws_struct_sq3_serialize(pss->pdb_artifact, - lsm_schema_sq3_map_artifact, - &o, (unsigned int)ap->uid)) { - lwsl_err("%s: failed artifact struct insert\n", - __func__); + ap = (sai_artifact_t *)pss->a.dest; - goto afail; - } + sai_task_uuid_to_event_uuid(event_uuid, ap->task_uuid); - /* - * recover the rowid - */ + /* + * Open the event-specific database object... the + * handle is closed when the stream closes, for whatever + * reason. + */ - lws_snprintf(s, sizeof(s), - "select rowid from artifacts " - "where timestamp=%llu", - (unsigned long long)ap->timestamp); + if (sais_event_db_ensure_open(pss->vhd, event_uuid, 0, + &pss->pdb_artifact)) { + lwsl_err("%s: unable to open event-specific " + "database\n", __func__); - if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, - sai_sql3_get_uint64_cb, &rid, NULL) != - SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, s, - sqlite3_errmsg(pss->pdb_artifact)); - goto afail; - } + lwsac_free(&pss->a.ac); + return -1; + } - /* - * Set the blob size on associated row - */ + /* + * Retreive the task object + */ - lws_snprintf(s, sizeof(s), - "update artifacts set blob=zeroblob(%llu) " - "where rowid=%llu", - (unsigned long long)ap->len, - (unsigned long long)rid); + lws_sql_purify(esc, ap->task_uuid, sizeof(esc)); + lws_snprintf(s, sizeof(s)," and uuid == \"%s\"", esc); + n = lws_struct_sq3_deserialize(pss->pdb_artifact, s, + NULL, lsm_schema_sq3_map_task, + &o, &ac, 0, 1); + if (n < 0 || !o.head) { + sais_event_db_close(vhd, &pss->pdb_artifact); + lwsl_notice("%s: no task of that id\n", __func__); + lwsac_free(&pss->a.ac); + return -1; + } - if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, - NULL, NULL, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, s, - sqlite3_errmsg(pss->pdb_artifact)); - goto afail; - } + task = (sai_task_t *)o.head; + n = strcmp(task->art_up_nonce, ap->artifact_up_nonce); - /* - * Open a blob on the associated row... the blob handle - * is closed when this stream closes for whatever - * reason. - */ + if (n) { + lwsl_err("%s: artifact nonce mismatch\n", + __func__); + goto afail; + } - if (sqlite3_blob_open(pss->pdb_artifact, "main", - "artifacts", "blob", (sqlite3_int64)rid, 1, - &pss->blob_artifact) != SQLITE_OK) { - lwsl_err("%s: unable to open blob\n", __func__); - goto afail; - } + /* + * The task the sender is sending us an artifact for + * exists. The sender knows the random upload nonce + * for that task's artifacts. + * + * Create a random download nonce unrelated to the + * random upload nonce (so knowing the download one + * won't let you upload anything). + * + * Create the artifact's entry in the event-specific + * database + */ - /* - * First time around, m == number of bytes let in buf - * after JSON, (bl - m) offset - */ - pss->bulk_binary_data = 1; - pss->artifact_length = ap->len; - } else { - m = (int)bl; - lwsl_info("%s: BUILDER_ARTIFACT: blob bulk\n", __func__); - } + sai_uuid16_create(pss->vhd->context, + ap->artifact_down_nonce); - if (m) { - lwsl_info("%s: blob write +%d, ofs %llu / %llu, len %d (0x%02x)\n", - __func__, (int)(bl - (unsigned int)m), - (unsigned long long)pss->artifact_offset, - (unsigned long long)pss->artifact_length, m, buf[0]); - if (sqlite3_blob_write(pss->blob_artifact, - (uint8_t *)buf + (bl - (unsigned int)m), (int)m, - (int)pss->artifact_offset)) { - lwsl_err("%s: writing blob failed\n", __func__); - goto afail; - } + lws_dll2_owner_clear(&o); + lws_dll2_add_head(&ap->list, &o); - lws_set_timeout(pss->wsi, PENDING_TIMEOUT_HTTP_CONTENT, 5); - pss->artifact_offset = pss->artifact_offset + (uint64_t)m; - } else - lwsl_info("%s: no m\n", __func__); + /* + * Create the task in event-specific database + */ + + if (lws_struct_sq3_serialize(pss->pdb_artifact, + lsm_schema_sq3_map_artifact, + &o, (unsigned int)ap->uid)) { + lwsl_err("%s: failed artifact struct insert\n", + __func__); - lwsl_info("%s: ofs %d, len %d\n", __func__, (int)pss->artifact_offset, (int)pss->artifact_length); + goto afail; + } - if (pss->artifact_offset == pss->artifact_length) { - int state; + /* + * recover the rowid + */ - lwsl_notice("%s: blob upload finished\n", __func__); - pss->bulk_binary_data = 0; + lws_snprintf(s, sizeof(s), + "select rowid from artifacts " + "where timestamp=%llu", + (unsigned long long)ap->timestamp); + + if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, + sai_sql3_get_uint64_cb, &rid, NULL) != + SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, s, + sqlite3_errmsg(pss->pdb_artifact)); + goto afail; + } - ap = (sai_artifact_t *)pss->a.dest; + /* + * Set the blob size on associated row + */ - lws_sql_purify(esc, ap->task_uuid, sizeof(esc)); - lws_snprintf(s, sizeof(s)," select state from tasks where uuid == \"%s\"", esc); - if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, - sql3_get_integer_cb, &state, NULL) != SQLITE_OK) { - lwsl_err("%s: %s: %s: fail\n", __func__, s, - sqlite3_errmsg(pss->pdb_artifact)); - goto bail; + lws_snprintf(s, sizeof(s), + "update artifacts set blob=zeroblob(%llu) " + "where rowid=%llu", + (unsigned long long)ap->len, + (unsigned long long)rid); + + if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, + NULL, NULL, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, s, + sqlite3_errmsg(pss->pdb_artifact)); + goto afail; + } + + /* + * Open a blob on the associated row... the blob handle + * is closed when this stream closes for whatever + * reason. + */ + + if (sqlite3_blob_open(pss->pdb_artifact, "main", + "artifacts", "blob", (sqlite3_int64)rid, 1, + &pss->blob_artifact) != SQLITE_OK) { + lwsl_err("%s: unable to open blob\n", __func__); + goto afail; + } + + /* + * First time around, m == number of bytes let in buf + * after JSON, (bl - m) offset + */ + pss->bulk_binary_data = 1; + pss->artifact_length = ap->len; + } else { + m = (int)bl; + lwsl_info("%s: BUILDER_ARTIFACT: blob bulk\n", __func__); } - sais_taskchange(pss->vhd->h_ss_websrv, ap->task_uuid, state); + if (m) { + lwsl_info("%s: blob write +%d, ofs %llu / %llu, len %d (0x%02x)\n", + __func__, (int)(bl - (unsigned int)m), + (unsigned long long)pss->artifact_offset, + (unsigned long long)pss->artifact_length, m, buf[0]); + if (sqlite3_blob_write(pss->blob_artifact, + (uint8_t *)buf + (bl - (unsigned int)m), (int)m, + (int)pss->artifact_offset)) { + lwsl_err("%s: writing blob failed\n", __func__); + goto afail; + } - goto afail; - } - break; + lws_set_timeout(pss->wsi, PENDING_TIMEOUT_HTTP_CONTENT, 5); + pss->artifact_offset = pss->artifact_offset + (uint64_t)m; + } else + lwsl_info("%s: no m\n", __func__); - case SAIM_WSSCH_BUILDER_RESOURCE_REQ: - res = (sai_resource_t *)pss->a.dest; + lwsl_info("%s: ofs %d, len %d\n", __func__, (int)pss->artifact_offset, (int)pss->artifact_length); - /* - * We get resource requests here, and also the handing back of - * assigned leases. The requests have the resname member and - * the lease yield messages don't. - */ + if (pss->artifact_offset == pss->artifact_length) { + int state; + + lwsl_notice("%s: blob upload finished\n", __func__); + pss->bulk_binary_data = 0; - if (!res->resname) { - sai_resource_requisition_t *rr; + ap = (sai_artifact_t *)pss->a.dest; + + lws_sql_purify(esc, ap->task_uuid, sizeof(esc)); + lws_snprintf(s, sizeof(s)," select state from tasks where uuid == \"%s\"", esc); + if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, + sql3_get_integer_cb, &state, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, s, + sqlite3_errmsg(pss->pdb_artifact)); + goto bail; + } + + sais_taskchange(pss->vhd->h_ss_websrv, ap->task_uuid, state); + + goto afail; + } + + m = 0; + + break; + + case SAIM_WSSCH_BUILDER_RESOURCE_REQ: + res = (sai_resource_t *)pss->a.dest; /* - * An assigned resource lease is being yielded + * We get resource requests here, and also the handing back of + * assigned leases. The requests have the resname member and + * the lease yield messages don't. */ - rr = sais_resource_lookup_lease_by_cookie(&vhd->server, - res->cookie); - if (!rr) { + if (!res->resname) { + sai_resource_requisition_t *rr; + + /* + * An assigned resource lease is being yielded + */ + + rr = sais_resource_lookup_lease_by_cookie(&vhd->server, + res->cookie); + if (!rr) { + /* + * He never got allocated... if he's on the + * queue delete him from there... if he doesn't + * exist on our side it's OK, just finish + */ + sais_resource_destroy_queued_by_cookie( + &vhd->server, res->cookie); + + return 0; + } + /* - * He never got allocated... if he's on the - * queue delete him from there... if he doesn't - * exist on our side it's OK, just finish + * Destroy the requisition, freeing any leased resources + * allocated to him */ - sais_resource_destroy_queued_by_cookie( - &vhd->server, res->cookie); + + sais_resource_rr_destroy(rr); return 0; } /* - * Destroy the requisition, freeing any leased resources - * allocated to him + * This is a new request for resources, find out the well-known + * resource to attach it to */ - sais_resource_rr_destroy(rr); - return 0; - } + wk = sais_resource_wellknown_by_name(&pss->vhd->server, + res->resname); + if (!wk) { + sai_resource_msg_t *mq; - /* - * This is a new request for resources, find out the well-known - * resource to attach it to - */ + /* + * Requested well-known resource doesn't exist + */ + + lwsl_info("%s: resource %s not well-known\n", __func__, + res->resname); + + mq = malloc(sizeof(*mq) + LWS_PRE + 256); + if (!mq) + return 0; + + memset(mq, 0, sizeof(*mq)); + + /* return with cookie but no amount == fail */ + mq->len = (size_t)lws_snprintf((char *)&mq[1] + LWS_PRE, 256, + "{\"schema\":\"com-warmcat-sai-resource\"," + "\"cookie\":\"%s\"}", res->cookie); + mq->msg = (char *)&mq[1] + LWS_PRE; - wk = sais_resource_wellknown_by_name(&pss->vhd->server, - res->resname); - if (!wk) { - sai_resource_msg_t *mq; + lws_dll2_add_tail(&mq->list, &pss->res_pending_reply_owner); + lws_callback_on_writable(pss->wsi); + + return 0; + } /* - * Requested well-known resource doesn't exist + * Create and queue the request on the right well-known + * resource manager, check if we can accept it */ - lwsl_info("%s: resource %s not well-known\n", __func__, - res->resname); - - mq = malloc(sizeof(*mq) + LWS_PRE + 256); - if (!mq) + rr = malloc(sizeof(*rr) + strlen(res->cookie) + 1); + if (!rr) return 0; + memset(rr, 0, sizeof(*rr)); + memcpy((char *)&rr[1], res->cookie, strlen(res->cookie) + 1); - memset(mq, 0, sizeof(*mq)); + rr->cookie = (char *)&rr[1]; + rr->lease_secs = res->lease; + rr->amount = res->amount; - /* return with cookie but no amount == fail */ + lws_dll2_add_tail(&rr->list_pss, &pss->res_owner); + lws_dll2_add_tail(&rr->list_resource_wellknown, &wk->owner); + lws_dll2_add_tail(&rr->list_resource_queued_leased, &wk->owner_queued); - mq->len = (size_t)lws_snprintf((char *)&mq[1] + LWS_PRE, 256, - "{\"schema\":\"com-warmcat-sai-resource\"," - "\"cookie\":\"%s\"}", res->cookie); - mq->msg = (char *)&mq[1] + LWS_PRE; + sais_resource_check_if_can_accept_queued(wk); + break; - lws_dll2_add_tail(&mq->list, &pss->res_pending_reply_owner); - lws_callback_on_writable(pss->wsi); + case SAIM_WSSCH_BUILDER_METRIC: + metric = (const sai_build_metric_t *)pss->a.dest; + sais_metrics_db_add(vhd, metric); - return 0; - } + { + uint8_t buf[2048]; + size_t used = 0; + lws_struct_serialize_t *js = lws_struct_json_serialize_create( + lsm_schema_build_metric, + LWS_ARRAY_SIZE(lsm_schema_build_metric), + 0, (void *)metric); + if (js) { + switch (lws_struct_json_serialize(js, buf, sizeof(buf), &used)) { + case LSJS_RESULT_CONTINUE: + assert(0); /* !!! we don't expect to generate anything that won't fit in one fragment */ + break; + case LSJS_RESULT_ERROR: + assert(0); /* we don't expect to not to be able to represent the metrics */ + break; + case LSJS_RESULT_FINISH: + memset(&info, 0, sizeof(info)); - /* - * Create and queue the request on the right well-known - * resource manager, check if we can accept it - */ + info.private_source_idx = SAI_WEBSRV_PB__PROXIED_FROM_BUILDER; + info.buf = (uint8_t *)buf; + info.len = used; + info.ss_flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; - rr = malloc(sizeof(*rr) + strlen(res->cookie) + 1); - if (!rr) - return 0; - memset(rr, 0, sizeof(*rr)); - memcpy((char *)&rr[1], res->cookie, strlen(res->cookie) + 1); - - rr->cookie = (char *)&rr[1]; - rr->lease_secs = res->lease; - rr->amount = res->amount; - - lws_dll2_add_tail(&rr->list_pss, &pss->res_owner); - lws_dll2_add_tail(&rr->list_resource_wellknown, &wk->owner); - lws_dll2_add_tail(&rr->list_resource_queued_leased, &wk->owner_queued); - - sais_resource_check_if_can_accept_queued(wk); - break; - - case SAIM_WSSCH_BUILDER_METRIC: - metric = (const sai_build_metric_t *)pss->a.dest; - sais_metrics_db_add(vhd, metric); - - { - uint8_t buf[2048]; - size_t used = 0; - lws_struct_serialize_t *js = lws_struct_json_serialize_create( - lsm_schema_build_metric, - LWS_ARRAY_SIZE(lsm_schema_build_metric), - 0, (void *)metric); - if (js) { - switch (lws_struct_json_serialize(js, buf, sizeof(buf), &used)) { - case LSJS_RESULT_CONTINUE: - assert(0); /* !!! we don't expect to generate anything that won't fit in one fragment */ - break; - case LSJS_RESULT_ERROR: - assert(0); /* we don't expect to not to be able to represent the metrics */ - break; - case LSJS_RESULT_FINISH: - sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, (const char *)buf, used, - SAI_WEBSRV_PB__PROXIED_FROM_BUILDER, LWSSS_FLAG_SOM | LWSSS_FLAG_EOM); - break; + if (sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info) < 0) + lwsl_warn("%s: unable to broadcast to web\n", __func__); + + break; + } + lws_struct_json_serialize_destroy(&js); } - lws_struct_json_serialize_destroy(&js); } + + lwsac_free(&pss->a.ac); + break; } - lwsac_free(&pss->a.ac); - break; - } + buf += ((int)bl - m); + bl = (size_t)m; + + } /* while (bl) */ return 0; @@ -1419,7 +1475,7 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, } - if (!pss->issue_task_owner.count || !pss->issue_task_owner.head) + if (!pss->issue_task_owner.head) return 0; /* nothing to send */ /*
Page fetched 0s ago, creation time: 10ms (vhost etag hits: 0%, cache hits: 0%)