Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / src / power / p-ws-server.c
Author[]Andy Green <andy@warmcat.com> 2025-08-28 18:54 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-08-31 13:22 UTC
Treec6e8f33d2506271b8be9c96a5388378f4e11330b   Raw Patch
 
scheduler: Remove fixed instance concept completely
scheduler: Remove fixed instance concept completely

Co-developed-by: Gemini 2.5 Pro
diff --git a/README.md b/README.md index 0829e7c..ab255b1 100644 --- a/README.md +++ b/README.md @@ -232,8 +232,9 @@ trees concurrently inside the platform. Tests have to take care to disambiguate which instance they are running on, since the network namespace is shared between instances that are running in the same sai-builder process on the same platform. An environment var -`SAI_INSTANCE_IDX` is available inside the each build context set to 0, 1, etc -according to the builder instance. +`SAI_INSTANCE_IDX` is available inside the each build context set to 0, 1, 33 etc +according to the builder instance. Similar to how fds are allocated in C, the +lowest unused number is reused each time something new is spawned by Sai. For network related tests, `SAI_INSTANCE_IDX` should be referred to when choosing, eg, a test server port so it will not conflict with what other diff --git a/src/builder/b-artifacts.c b/src/builder/b-artifacts.c index 17e5ce2..f71de4b 100644 --- a/src/builder/b-artifacts.c +++ b/src/builder/b-artifacts.c @@ -157,6 +157,12 @@ saib_artifact_state(void *userobj, void *sh, lws_ss_constate_t state, ap->fd = -1; } unlink(ap->path); + + ap->ns->count_artifacts--; + if (!ap->ns->count_artifacts) { + lwsl_notice("%s: last artifact completed, destroying ns now\n", __func__); + saib_task_destroy(ap->ns); + } break; case LWSSSCS_CONNECTED: diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index 7b21996..62ccd72 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -253,6 +253,8 @@ send_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). * * That requires statefully rotating through... * @@ -367,10 +369,14 @@ send_logs: /* * 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: drained and empty\n", __func__); + lwsl_notice("%s: logs cache drained and empty\n", __func__); ns->finished_when_logs_drained = 0; - saib_set_ns_state(ns, NSSTATE_UPLOADING_ARTIFACTS); + if (ns->state != NSSTATE_FAILED) + saib_set_ns_state(ns, NSSTATE_UPLOADING_ARTIFACTS); } break; @@ -407,6 +413,7 @@ cleanup_on_ss_destroy(struct lws_dll2 *d, void *user) lws_container_of(d, struct sai_nspawn, list); if (ns->spm == spm) { + lwsl_warn("%s: ns->spm %p, spm %p\n", __func__, ns->spm, spm); /* * This pss is about to go away, make sure the ns * can't reference it any more no matter what happens @@ -590,7 +597,7 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, break; case LWSSSCS_CONNECTED: - lwsl_user("%s: CONNECTED: %p\n", __func__, spm->ss); + lwsl_ss_user(spm->ss, "CONNECTED"); spm->phase = PHASE_START_ATTACH; /* Initialize the load report SUL timer for this server connection */ lws_sul_cancel(&spm->sul_load_report); @@ -602,7 +609,7 @@ saib_m_state(void *userobj, void *sh, lws_ss_constate_t state, * clean up any ongoing spawns related to this connection */ - lwsl_user("%s: DISCONNECTED\n", __func__); + lwsl_ss_user(spm->ss, "DISCONNECTED"); lws_sul_cancel(&spm->sul_load_report); lws_dll2_foreach_safe(&builder.sai_plat_owner, spm, cleanup_on_ss_disconnect); diff --git a/src/builder/b-metrics.c b/src/builder/b-metrics.c index 170fcdc..5532cfb 100644 --- a/src/builder/b-metrics.c +++ b/src/builder/b-metrics.c @@ -68,6 +68,75 @@ saib_get_free_ram_kib(void) statex.dwLength = sizeof(statex); GlobalMemoryStatusEx(&statex); return (unsigned int)(statex.ullAvailPhys / 1024); +#elif defined(__APPLE__) + int mib[2]; + size_t len; + uint64_t total_mem; + + mib[0] = CTL_HW; + mib[1] = HW_MEMSIZE; + len = sizeof(total_mem); + sysctl(mib, 2, &total_mem, &len, NULL, 0); + + return (unsigned int)(total_mem / 1024); +#elif defined(_WIN32) + MEMORYSTATUSEX statex; + statex.dwLength = sizeof(statex); + GlobalMemoryStatusEx(&statex); + return (unsigned int)(statex.ullTotalPhys / 1024); +#else + return 0; +#endif +} + +unsigned int +saib_get_total_ram_kib(void) +{ +#if defined(__linux__) + char buf[256]; + FILE *f; + unsigned int total_kib = 0; + + f = fopen("/proc/meminfo", "r"); + if (!f) + return 0; + + while (fgets(buf, sizeof(buf), f)) { + if (sscanf(buf, "MemTotal: %u kB", &total_kib) == 1) + break; + } + + fclose(f); + return total_kib; +#else + return 0; +#endif +} + +unsigned int +saib_get_total_disk_kib(const char *path) +{ +#if defined(__linux__) + struct statvfs s; + + if (statvfs(path, &s)) + return 0; + + return (unsigned int)((uint64_t)s.f_blocks * s.f_frsize / 1024); +#elif defined(__APPLE__) + struct statfs s; + + if (statfs(path, &s)) + return 0; + + return (unsigned int)((uint64_t)s.f_blocks * (uint64_t)s.f_bsize / 1024); +#elif defined(_WIN32) + ULARGE_INTEGER total_bytes; + + if (!GetDiskFreeSpaceExA(path, NULL, &total_bytes, NULL)) + return 0; + + return (unsigned int)(total_bytes.QuadPart / 1024); #else return 0; #endif diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index 1c3a022..6935808 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -89,7 +89,7 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, switch (reason) { case LWS_CALLBACK_RAW_CLOSE_FILE: - lwsl_warn("%s: RAW_CLOSE_FILE at %llu, wsi %p: fd: %d, stdfd: %d\n", + lwsl_info("%s: RAW_CLOSE_FILE at %llu, wsi %p: fd: %d, stdfd: %d\n", __func__, (unsigned long long)lws_now_usecs(), wsi, lws_get_socket_fd(wsi), lws_spawn_get_stdfd(wsi)); @@ -102,7 +102,6 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, __func__); } - lwsl_wsi_err(wsi, "CLOSING: op %p, op->lsp %p", op, op ? op->lsp : NULL); if (op && op->lsp) { lws_spawn_stdwsi_closed(op->lsp, wsi); if (ns) @@ -150,6 +149,11 @@ callback_sai_stdwsi(struct lws *wsi, enum lws_callback_reasons reason, struct lws_protocols protocol_stdxxx = { "sai-stdxxx", callback_sai_stdwsi, 0, 0 }; +/* + * We are called when the process completed and has been reaped at + * lsp level, and we know that all the stdwsi related to the process + * are closed. + */ static void sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, @@ -249,7 +253,14 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, } if (op->spawn) { - sai_build_metric_t *m = malloc(sizeof(*m)); + sai_build_metric_t *m; + + if (!ns->spm) { + lwsl_err("%s: NULL ns->spm", __func__); + goto skip; + } + + m = malloc(sizeof(*m)); if (m) { char hash_input[8192]; @@ -288,6 +299,7 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, } } +skip: ns->current_step++; /* step succeeded, wait for next instruction */ @@ -296,9 +308,6 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, if (op->spawn) free(op->spawn); - saib_task_grace(ns); - saib_set_ns_state(ns, NSSTATE_DONE); - /* * add a final zero-length log with the retcode to the list of pending * logs @@ -310,6 +319,13 @@ sai_lsp_reap_cb(void *opaque, const lws_spawn_resource_us_t *res, siginfo_t *si, __func__, ns->chunk_cache.count, ns->spm ? ns->spm->logs_in_flight : -99); + /* + * saib_task_grace(ns) sets ns->finished_when_logs_drained + */ + + saib_task_grace(ns); + saib_set_ns_state(ns, NSSTATE_DONE); + if (ns) ns->op = NULL; free(op); @@ -320,6 +336,9 @@ 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); @@ -492,7 +511,7 @@ saib_spawn_script(struct sai_nspawn *ns) #if defined(WIN32) n = lws_snprintf(st, sizeof(st), ns->current_step ? runscript_win_next : runscript_win_first, - ns->instance_idx, + ns->instance_ordinal, 1, respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, builder.home, @@ -510,7 +529,7 @@ saib_spawn_script(struct sai_nspawn *ns) n = lws_snprintf(st, sizeof(st), script_template, builder.home, ns->fsm.ovname, - ns->project_name, ns->ref, ns->instance_idx, + ns->project_name, ns->ref, ns->instance_ordinal, 1, respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, @@ -530,7 +549,7 @@ saib_spawn_script(struct sai_nspawn *ns) cmd[0] = ns->script_path; #if defined(__linux__) - lws_snprintf(cgroup, sizeof(cgroup), "inst-%u-%d", (unsigned int)getpid(), ns->instance_idx); + lws_snprintf(cgroup, sizeof(cgroup), "inst-%s", ns->task->uuid); #endif memset(&info, 0, sizeof(info)); diff --git a/src/builder/b-private.h b/src/builder/b-private.h index 5de8992..6d14267 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -233,4 +233,14 @@ unsigned int saib_get_free_ram_kib(void); unsigned int +saib_get_total_ram_kib(void); + +unsigned int saib_get_free_disk_kib(const char *path); + +unsigned int +saib_get_total_disk_kib(const char *path); + +int +saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp, struct lws_vhost **); + diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index a2ef7bd..28a478e 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -197,10 +197,17 @@ pvo_resproxy = { /* starting point for resproxy */ NULL, /* "child" pvo linked-list */ "protocol-resproxy", /* protocol name we belong to on this vhost */ "ok" /* ignored */ - };; + }; -static int -saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp) +void * +saib_thread_suspend(void *d) +{ + return NULL; +} + +int +saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp, + struct lws_vhost **vhost) { struct lws_context_creation_info info; @@ -223,7 +230,8 @@ saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp) lwsl_notice("%s: %s.%s\n", __func__, info.vhost_name, lp->sockpath); - if (!lws_create_vhost(context, &info)) { + *vhost = lws_create_vhost(context, &info); + if (!*vhost) { lwsl_notice("%s: failed to create vh %s\n", __func__, info.vhost_name); return -1; @@ -400,9 +408,6 @@ static int app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link, int current, int target) { - struct lws_context *context = lws_system_context_from_system_mgr(mgr); - char pur[128]; - /* * For the things we care about, let's notice if we are trying to get * past them when we haven't solved them yet, and make the system @@ -432,85 +437,6 @@ app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link, * For each platform... */ - lws_start_foreach_dll_safe(struct lws_dll2 *, mp, mp1, - builder.sai_plat_owner.head) { - struct sai_plat *sp = lws_container_of(mp, struct sai_plat, - sai_plat_list); - -#if defined(WIN32) - sp->windows = 1; -#endif - - /* - * ... for each nspawn on the platform... - */ - - lws_start_foreach_dll_safe(struct lws_dll2 *, np, np1, - sp->nspawn_owner.head) { - struct sai_nspawn *ns = - lws_container_of(np, struct sai_nspawn, list); - char *p; - int n; - - lws_strncpy(pur, sp->name, sizeof(pur)); - lws_filename_purify_inplace(pur); - p = pur; - while ((p = strchr(p, '/'))) - *p = '_'; - - /* - * Proxy the control logging channel (this - * is the one that has sai progress info) - */ - - lws_snprintf(ns->slp_control.sockpath, - sizeof(ns->slp_control.sockpath), -#if defined(__linux__) - UDS_PATHNAME_LOGPROXY".%s.%d.saib", -#else - UDS_PATHNAME_LOGPROXY"/%s.%d.saib", -#endif - pur, ns->instance_idx); - - ns->slp_control.ns = ns; - ns->slp_control.log_channel_idx = 3; - - if (saib_create_listen_uds(context, &ns->slp_control)) { - lwsl_err("%s: Failed to create ctl log proxy listen UDS %s\n", - __func__, ns->slp_control.sockpath); - return -1; - } - - /* - * For each additional logging channel... - */ - - for (n = 0; n < (int)LWS_ARRAY_SIZE(ns->slp); n++) { - - /* - * ... create a UDS listening proxy - */ - - lws_snprintf(ns->slp[n].sockpath, - sizeof(ns->slp[n].sockpath), -#if defined(__linux__) - UDS_PATHNAME_LOGPROXY".%s.%d.tty%d", -#else - UDS_PATHNAME_LOGPROXY"/%s.%d.tty%d", -#endif - pur, ns->instance_idx, n); - - ns->slp[n].ns = ns; - ns->slp[n].log_channel_idx = n + 4; - - if (saib_create_listen_uds(context, &ns->slp[n])) { - lwsl_err("%s: Failed to create log proxy listen UDS %s\n", - __func__, ns->slp[n].sockpath); - return -1; - } - } - } lws_end_foreach_dll_safe(np, np1); - } lws_end_foreach_dll_safe(mp, mp1); /* * Create the resource proxy listeners, one per server link diff --git a/src/builder/b-task.c b/src/builder/b-task.c index a4c4448..b90f080 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -27,6 +27,40 @@ #include "b-private.h" +static int +saib_can_accept_task(sai_task_t *task) +{ +#if 0 + unsigned int free_ram = saib_get_free_ram_kib(); + unsigned int total_ram = saib_get_total_ram_kib(); + unsigned int free_disk = saib_get_free_disk_kib(builder.home); + unsigned int total_disk = saib_get_total_disk_kib(builder.home); +// int cpu_load = saib_get_system_cpu(&builder); + + if (total_ram && + (free_ram - task->est_peak_mem_kib) < (total_ram / 10) * 3) { + lwsl_notice("%s: reject task %s: not enough RAM\n", __func__, + task->uuid); + return 1; + } + + if (total_disk && + (free_disk - task->est_disk_kib) < (total_disk / 10) * 2) { + lwsl_notice("%s: reject task %s: not enough disk space\n", + __func__, task->uuid); + return 1; + } + +/* if (cpu_load >= 0 && (cpu_load + (int)task->est_cpu_load_pct) > 50) { + lwsl_notice("%s: reject task %s: CPU load too high\n", + __func__, task->uuid); + return 1; + } +*/ +#endif + return 0; /* acceptable */ +} + #if !defined(WIN32) static char csep = '/'; #else @@ -112,7 +146,6 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, return 0; rej = malloc(sizeof(*rej)); - if (!rej) return -1; @@ -142,6 +175,8 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, void saib_task_destroy(struct sai_nspawn *ns) { + int n; + ns->finished_when_logs_drained = 0; lws_sul_cancel(&ns->sul_cleaner); lws_sul_cancel(&ns->sul_task_cancel); @@ -155,6 +190,13 @@ saib_task_destroy(struct sai_nspawn *ns) ns->spm->phase = PHASE_START_ATTACH; if (lws_ss_request_tx(ns->spm->ss)) return; + + /* + * If spm is holding on to us as the last reference point, + * we can't be used any more since we are goneski + */ + if (ns->spm->last_logging_nspawn == &ns->list) + ns->spm->last_logging_nspawn = NULL; } if (ns->list.owner && ns->list.owner->count == 1) { @@ -196,6 +238,23 @@ saib_task_destroy(struct sai_nspawn *ns) if (ns->script_path[0]) unlink(ns->script_path); + for (n = 0; n < (int)LWS_ARRAY_SIZE(ns->vhosts); n++) + if (ns->vhosts[n]) { + lws_vhost_destroy(ns->vhosts[n]); + ns->vhosts[n] = NULL; + } + + if (ns->slp_control.sockpath[0]) + unlink(ns->slp_control.sockpath); + for (n = 0; n < (int)LWS_ARRAY_SIZE(ns->slp); n++) + if (ns->slp[n].sockpath[0]) + unlink(ns->slp[n].sockpath); + + /* + * If stdwsi are lurking around, we can't destroy the ns, + * since they will touch it during their close handling. + */ + lws_dll2_remove(&ns->list); free(ns); } @@ -205,12 +264,20 @@ saib_sub_cleaner_cb(lws_sorted_usec_list_t *sul) { struct sai_nspawn *ns = lws_container_of(sul, struct sai_nspawn, sul_cleaner); - lwsl_warn("%s: .%d: Destroying task after grace period\n", - __func__, ns->instance_idx); + lwsl_warn("%s: Task completion grace period ended with ns alive\n", __func__); saib_task_destroy(ns); } +void +saib_task_grace(struct sai_nspawn *ns) +{ + 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); +} + + static int artifact_glob_cb(void *data, const char *path) { @@ -218,7 +285,7 @@ artifact_glob_cb(void *data, const char *path) const char *p, *ph = NULL; struct lws_ss_handle *h; char upp[256], s[384]; - sai_artifact_t *ap; + sai_artifact_t *ap = NULL; int n; /* @@ -271,15 +338,15 @@ artifact_glob_cb(void *data, const char *path) /* take a copy so we can unlink the path later */ lws_strncpy(ap->path, upp, sizeof(ap->path)); + lwsl_notice("%s: artifact ss created '%s'\n", __func__, ap->path); + ns->count_artifacts++; + lws_strncpy(ap->task_uuid, ns->task->uuid, sizeof(ap->task_uuid)); lws_strncpy(ap->artifact_up_nonce, ns->task->art_up_nonce, sizeof(ap->artifact_up_nonce)); lws_strncpy(ap->blob_filename, ph, sizeof(ap->blob_filename)); ap->timestamp = (uint64_t)lws_now_usecs(); - lwsl_notice("%s: artifact ss created '%s'\n", __func__, ap->path); - - /* * We need to set the metadata items for the post urlargs. spm->url is * something like "wss://warmcat.com/sai/builder"... we will send JSON @@ -294,7 +361,8 @@ artifact_glob_cb(void *data, const char *path) } /* - * We're finished with the nspawn / task one way or the other, but there's + * We're finished with the nspawn / task one way or the other, specifically + * all the stdwsi are closed and we reaped the lws_spawn_piped, but there's * still stuff we need to send out. Give it some time then force destruction * of the task and reset the nspawn. */ @@ -399,14 +467,14 @@ scan: if (ts.e == LWS_TOKZE_ENDED) break; } -} -void -saib_task_grace(struct sai_nspawn *ns) -{ - ns->finished_when_logs_drained = 1; - lws_sul_schedule(builder.context, 0, &ns->sul_cleaner, - saib_sub_cleaner_cb, 20 * LWS_USEC_PER_SEC); + if (!ns->count_artifacts) { + lwsl_notice("%s: no artifacts, destroying ns now\n", __func__); + /* no artifacts to hang around for... nuke the ns now */ + lws_sul_schedule(builder.context, 0, &ns->sul_cleaner, + saib_sub_cleaner_cb, 1); + } else + lwsl_notice("%s: created / waiting on %d artifact uploads\n", __func__, ns->count_artifacts); } static void @@ -503,6 +571,22 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) } /* + * There's not already an existing step we accepted, + * using the same uuid? + */ + + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, sp->nspawn_owner.head) { + struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list); + + if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { + lwsl_err("%s: server offered task that's already running\n", __func__); + return 0; + } + } lws_end_foreach_dll_safe(d, d1); + + + + /* * store a copy of the toplevel ac used for the deserialization * into the outer part of the c builder wrapper */ @@ -520,10 +604,8 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) n = 0; ns = NULL; - lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, - sp->nspawn_owner.head) { - struct sai_nspawn *xns = lws_container_of(d, - struct sai_nspawn, list); + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, sp->nspawn_owner.head) { + struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list); if (xns->task && !strcmp(xns->task->uuid, task->uuid)) { ns = xns; break; @@ -532,14 +614,91 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) ns = xns; } lws_end_foreach_dll_safe(d, d1); + if (saib_can_accept_task(task)) { /* not accepted */ + if (saib_queue_task_status_update(sp, spm, task->uuid)) + return -1; + return 0; + } + if (!ns) { + char pur[128], *p, ordinal_acc[SAI_BUILDER_INSTANCE_LIMIT]; + int n; + ns = malloc(sizeof(*ns)); if (!ns) return -1; memset(ns, 0, sizeof(*ns)); ns->builder = &builder; ns->sp = sp; + + /* + * Find the lowest free ordinal and use that. It doesn't + * have any meaning for us, but the project being built needs + * it in SAI_INSTANCE_IDX so ctest can use, eg, test ports + * that don't conflict with any other running instance. + */ + + memset(ordinal_acc, 0, sizeof(ordinal_acc)); + lws_start_foreach_dll_safe(struct lws_dll2 *, d, d1, sp->nspawn_owner.head) { + struct sai_nspawn *xns = lws_container_of(d, struct sai_nspawn, list); + assert(xns->instance_ordinal < (int)sizeof(ordinal_acc)); + ordinal_acc[xns->instance_ordinal] = 1; + } lws_end_foreach_dll_safe(d, d1); + + for (n = 0; n < (int)sizeof(ordinal_acc); n++) + if (ordinal_acc[n] == 0) { + ns->instance_ordinal = n; + break; + } + + lws_dll2_add_tail(&ns->list, &sp->nspawn_owner); + + if (strstr(task->script, "sai-device")) { + lws_strncpy(pur, sp->name, sizeof(pur)); + lws_filename_purify_inplace(pur); + p = pur; + while ((p = strchr(p, '/'))) + *p = '_'; + + lws_snprintf(ns->slp_control.sockpath, + sizeof(ns->slp_control.sockpath), +#if defined(__linux__) + UDS_PATHNAME_LOGPROXY".%s.saib", +#else + UDS_PATHNAME_LOGPROXY"/%s.saib", +#endif + task->uuid); + ns->slp_control.ns = ns; + ns->slp_control.log_channel_idx = 3; + if (saib_create_listen_uds(builder.context, &ns->slp_control, + &ns->vhosts[0])) { + lwsl_err("%s: Failed to create ctl log proxy listen UDS %s\n", + __func__, ns->slp_control.sockpath); + return -1; + } + + for (n = 0; n < (int)LWS_ARRAY_SIZE(ns->slp); n++) { + lws_snprintf(ns->slp[n].sockpath, + sizeof(ns->slp[n].sockpath), +#if defined(__linux__) + UDS_PATHNAME_LOGPROXY".%s.tty%d", +#else + UDS_PATHNAME_LOGPROXY"/%s.tty%d", +#endif + task->uuid, n); + + ns->slp[n].ns = ns; + ns->slp[n].log_channel_idx = n + 4; + + if (saib_create_listen_uds(builder.context, &ns->slp[n], + &ns->vhosts[n + 1])) { + lwsl_err("%s: Failed to create log proxy listen UDS %s\n", + __func__, ns->slp[n].sockpath); + return -1; + } + } + } } // lwsl_hexdump_warn(task->build, strlen(task->build)); diff --git a/src/common/include/private.h b/src/common/include/private.h index 4219c94..971ac19 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -35,6 +35,8 @@ #define UDS_PATHNAME_RESPROXY "/var/run/com.warmcat.com.saib.resproxy" #endif +#define SAI_BUILDER_INSTANCE_LIMIT 256 + struct sai_plat; struct sai_builder; struct saib_opaque_spawn; @@ -122,6 +124,13 @@ typedef struct { int state; int uid; int build_step; + + /* estimations for builder resource consumption */ + unsigned int est_peak_mem_kib; + unsigned int est_cpu_load_pct; + unsigned int est_disk_kib; + + char told_ongoing; } sai_task_t; typedef struct sai_plat sai_plat_t; @@ -148,6 +157,7 @@ struct sai_nspawn { /* convenient place to store it */ struct saib_logproxy slp_control; struct saib_logproxy slp[2]; + struct lws_vhost *vhosts[3]; lws_dll2_t list; /* sai_plat owner lists sai_nspawns */ struct sai_builder *builder; @@ -155,8 +165,6 @@ struct sai_nspawn { struct saib_opaque_spawn *op; sai_task_t *task; - lws_dll2_owner_t artifact_owner; /* struct artifact_path */ - lws_spawn_resource_us_t res; lws_sorted_usec_list_t sul_cleaner; @@ -184,16 +192,16 @@ struct sai_nspawn { const char *git_repo_url; int retcode; - int instance_idx; + int instance_ordinal; + int count_artifacts; + int current_step; + int build_step_count; uint8_t spins; uint8_t state; /* NSSTATE_ */ uint8_t stdcount; uint8_t term_budget; - int current_step; - int build_step_count; - uint8_t finished_when_logs_drained:1; uint8_t state_changed:1; uint8_t user_cancel:1; @@ -271,6 +279,7 @@ typedef struct { struct lws_ss_handle *ss; void *opaque_data; + struct sai_nspawn *ns; char task_uuid[65]; char artifact_up_nonce[33]; char artifact_down_nonce[33]; @@ -509,7 +518,7 @@ extern const lws_struct_map_t lsm_schema_map_ta[1], lsm_schema_map_plat_simple[1], lsm_event[10], - lsm_task[23], + lsm_task[26], lsm_log[7], lsm_artifact[8], lsm_plat_list[1], diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 4453ddf..a36df7a 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -163,6 +163,9 @@ const lws_struct_map_t lsm_task[] = { LSM_STRING_PTR (sai_task_t, git_repo_url, "git_repo_url"), LSM_CARRAY (sai_task_t, script, "script"), LSM_SIGNED (sai_task_t, build_step, "build_step"), + LSM_UNSIGNED (sai_task_t, est_peak_mem_kib, "est_peak_mem_kib"), + LSM_UNSIGNED (sai_task_t, est_cpu_load_pct, "est_cpu_load_pct"), + LSM_UNSIGNED (sai_task_t, est_disk_kib, "est_disk_kib"), }; const lws_struct_map_t lsm_schema_json_map_task[] = { diff --git a/src/power/p-intake.c b/src/power/p-intake.c index bb2d386..01a3ef7 100644 --- a/src/power/p-intake.c +++ b/src/power/p-intake.c @@ -35,7 +35,7 @@ static int -callback_ws_power(struct lws *wsi, enum lws_callback_reasons reason, void *user, +p_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, void *in, size_t len) { struct vhd *vhd = (struct vhd *)lws_protocol_vh_priv_get( @@ -100,7 +100,7 @@ callback_ws_power(struct lws *wsi, enum lws_callback_reasons reason, void *user, case LWS_CALLBACK_CLOSED: lwsac_free(&pss->query_ac); - lwsl_user("%s: CLOSED builder conn\n", __func__); + lwsl_wsi_user(wsi, "CLOSED builder->power connection", __func__); /* remove pss from vhd->builders */ lws_dll2_remove(&pss->same); @@ -169,4 +169,4 @@ passthru: } const struct lws_protocols protocol_ws_power = - { "com-warmcat-sai-power", callback_ws_power, sizeof(struct pss), 0 }; + { "com-warmcat-sai-power", p_callback_ws, sizeof(struct pss), 0 }; diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 6dc669d..573c3e4 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -386,7 +386,7 @@ sai_get_head_status(struct vhd *vhd, const char *projname) } static int -callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, +s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, void *in, size_t len) { struct vhd *vhd = (struct vhd *)lws_protocol_vh_priv_get( @@ -759,7 +759,7 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, case LWS_CALLBACK_CLOSED: lwsac_free(&pss->query_ac); - lwsl_user("%s: CLOSED builder conn\n", __func__); + lwsl_wsi_user(wsi, "sai-server: CLOSED sai-web conn"); /* remove pss from vhd->builders (active connection list) */ lws_dll2_remove(&pss->same); @@ -867,4 +867,4 @@ passthru: } const struct lws_protocols protocol_ws = - { "com-warmcat-sai", callback_ws, sizeof(struct pss), 0 }; + { "com-warmcat-sai", s_callback_ws, sizeof(struct pss), 0 }; diff --git a/src/server/s-task.c b/src/server/s-task.c index 21525d3..f50986f 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -82,6 +82,10 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, sai_event_t *e = NULL; lws_dll2_owner_t o; int n, task_ostate; + int ostate = state; + + if (state == SAIES_STEP_SUCCESS) + state = SAIES_BEING_BUILT; /* * Extract the event uuid from the task uuid @@ -273,11 +277,8 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, sais_event_db_close(vhd, (sqlite3 **)&e->pdb); lwsac_free(&ac); - if (state == SAIES_STEP_SUCCESS) { + if (ostate == SAIES_STEP_SUCCESS) sais_continue_task(vhd, task_uuid); - state = SAIES_BEING_BUILT; - } - return 0; @@ -679,6 +680,41 @@ bail: return 1; } +static void +sais_get_task_metrics_estimates(struct vhd *vhd, sai_task_t *task) +{ + char query[256]; + sqlite3_stmt *stmt; + + task->est_peak_mem_kib = 256 * 1024; /* 256MiB default */ + task->est_cpu_load_pct = 10; + task->est_disk_kib = 1024 * 1024; /* 1GiB default */ + + if (!vhd->pdb_metrics) + return; + + lws_snprintf(query, sizeof(query), + "SELECT AVG(peak_mem_rss), AVG(us_cpu_user), " + "AVG(stg_bytes), AVG(wallclock_us) " + "FROM build_metrics WHERE key = '%s'", + task->taskname); + + if (sqlite3_prepare_v2(vhd->pdb_metrics, query, -1, &stmt, NULL) != SQLITE_OK) + return; + + if (sqlite3_step(stmt) == SQLITE_ROW) { + uint64_t avg_us_cpu = (uint64_t)sqlite3_column_int64(stmt, 1); + uint64_t avg_wallclock = (uint64_t)sqlite3_column_int64(stmt, 3); + + task->est_peak_mem_kib = (unsigned int)(sqlite3_column_int(stmt, 0) / 1024); + if (avg_wallclock) + task->est_cpu_load_pct = (unsigned int)((avg_us_cpu * 100) / avg_wallclock); + task->est_disk_kib = (unsigned int)(sqlite3_column_int(stmt, 2) / 1024); + } + + sqlite3_finalize(stmt); +} + int sais_task_cancel(struct vhd *vhd, const char *task_uuid) { @@ -718,6 +754,36 @@ sais_task_cancel(struct vhd *vhd, const char *task_uuid) return 0; } +static int +sais_task_stop_on_builders(struct vhd *vhd, const char *task_uuid) +{ + sai_cancel_t *can; + + lwsl_notice("%s: builders count %d\n", __func__, vhd->builders.count); + + /* + * For every pss that we have from builders... + */ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->builders.head) { + struct pss *pss = lws_container_of(p, struct pss, same); + + /* + * ... queue the task cancel message + */ + can = malloc(sizeof *can); + if (!can) + return -1; + memset(can, 0, sizeof(*can)); + + lws_strncpy(can->task_uuid, task_uuid, sizeof(can->task_uuid)); + lws_dll2_add_tail(&can->list, &pss->task_cancel_owner); + lws_callback_on_writable(pss->wsi); + + } lws_end_foreach_dll(p); + + return 0; +} + /* * Keep the task record itself, but remove all logs and artifacts related to * it and reset the task state back to WAITING. @@ -730,6 +796,8 @@ sais_task_reset(struct vhd *vhd, const char *task_uuid) sqlite3 *pdb = NULL; int ret; + lwsl_notice("%s: task reset %s\n", __func__, task_uuid); + if (!task_uuid[0]) return SAI_DB_RESULT_OK; @@ -774,7 +842,7 @@ sais_task_reset(struct vhd *vhd, const char *task_uuid) sais_set_task_state(vhd, NULL, NULL, task_uuid, SAIES_WAITING, 1, 1); - sais_task_cancel(vhd, task_uuid); + sais_task_stop_on_builders(vhd, task_uuid); /* * Reassess now if there's a builder we can match to a pending task @@ -788,6 +856,8 @@ sais_task_reset(struct vhd *vhd, const char *task_uuid) */ sais_platforms_with_tasks_pending(vhd); + lwsl_notice("%s: exiting OK\n", __func__); + return SAI_DB_RESULT_OK; } @@ -960,6 +1030,8 @@ sais_continue_task(struct vhd *vhd, const char *task_uuid) *task = *task_template; lwsac_free(&ac); + sais_get_task_metrics_estimates(vhd, task); + build_step = task->build_step; /* get the event */ @@ -976,6 +1048,7 @@ sais_continue_task(struct vhd *vhd, const char *task_uuid) event = lws_container_of(o_event.head, sai_event_t, list); task->one_event = event; + task->repo_name = event->repo_name; task->git_ref = event->ref; task->git_hash = event->hash; diff --git a/src/server/s-websrv.c b/src/server/s-websrv.c index 789898a..e96f50e 100644 --- a/src/server/s-websrv.c +++ b/src/server/s-websrv.c @@ -49,9 +49,6 @@ typedef struct sai_sul_retry_ctx { uint8_t op; /* SAIS_WS_WEBSRV_RX_... */ } sai_sul_retry_ctx_t; -static void -sais_websrv_retry_cb(lws_sorted_usec_list_t *sul); - typedef struct websrvss_srv { struct lws_ss_handle *ss; struct vhd *vhd; @@ -557,39 +554,6 @@ sais_plat_reset(struct vhd *vhd, const char *event_uuid, const char *platform) return SAI_DB_RESULT_OK; } -static void -sais_websrv_retry_cb(lws_sorted_usec_list_t *sul) -{ - sai_sul_retry_ctx_t *ctx = lws_container_of(sul, sai_sul_retry_ctx_t, sul); - sai_db_result_t r = SAI_DB_RESULT_ERROR; - - switch(ctx->op) { - case SAIS_WS_WEBSRV_RX_TASKRESET: - r = sais_task_reset(ctx->vhd, ctx->uuid); - break; - case SAIS_WS_WEBSRV_RX_EVENTRESET: - r = sais_event_reset(ctx->vhd, ctx->uuid); - break; - case SAIS_WS_WEBSRV_RX_EVENTDELETE: - r = sais_event_delete(ctx->vhd, ctx->uuid); - break; - case SAIS_WS_WEBSRV_RX_PLATRESET: - r = sais_plat_reset(ctx->vhd, ctx->uuid, ctx->platform); - break; - } - - if (r == SAI_DB_RESULT_BUSY && ctx->retries-- > 0) { - lwsl_notice("Retrying op %d for %s\n", ctx->op, ctx->uuid); - lws_sul_schedule(ctx->vhd->context, 0, &ctx->sul, - sais_websrv_retry_cb, 250 * LWS_US_PER_MS); - return; - } - - if (r != SAI_DB_RESULT_OK) - lwsl_err("Failed op %d for %s after retries\n", ctx->op, ctx->uuid); - - free(ctx); -} static lws_ss_state_return_t websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) @@ -622,26 +586,13 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) switch (a.top_schema_index) { case SAIS_WS_WEBSRV_RX_TASKRESET: { - sai_db_result_t r; - ei = (sai_browse_rx_evinfo_t *)a.dest; if (sais_validate_id(ei->event_hash, SAI_TASKID_LEN)) goto soft_error; - r = sais_task_reset(m->vhd, ei->event_hash); - if (r == SAI_DB_RESULT_BUSY) { - sai_sul_retry_ctx_t *ctx = malloc(sizeof(*ctx)); - - if (!ctx) - break; - - ctx->vhd = m->vhd; - lws_strncpy(ctx->uuid, ei->event_hash, sizeof(ctx->uuid)); - ctx->retries = 10; - ctx->op = SAIS_WS_WEBSRV_RX_TASKRESET; - lws_sul_schedule(m->vhd->context, 0, &ctx->sul, - sais_websrv_retry_cb, 250 * LWS_US_PER_MS); - } + lwsl_ss_warn(m->ss, "SAIS_WS_WEBSRV_RX_TASKRESET: %s: received", ei->event_hash); + if (sais_task_reset(m->vhd, ei->event_hash)) + lwsl_ss_err(m->ss, "taskreset failed"); break; } @@ -650,23 +601,14 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) sai_db_result_t r; ei = (sai_browse_rx_evinfo_t *)a.dest; + if (sais_validate_id(ei->event_hash, SAI_EVENTID_LEN)) goto soft_error; r = sais_event_reset(m->vhd, ei->event_hash); - if (r == SAI_DB_RESULT_BUSY) { - sai_sul_retry_ctx_t *ctx = malloc(sizeof(*ctx)); + if (r) + lwsl_ss_err(m->ss, "eventreset failed"); - if (!ctx) - break; - - ctx->vhd = m->vhd; - lws_strncpy(ctx->uuid, ei->event_hash, sizeof(ctx->uuid)); - ctx->retries = 10; - ctx->op = SAIS_WS_WEBSRV_RX_EVENTRESET; - lws_sul_schedule(m->vhd->context, 0, &ctx->sul, - sais_websrv_retry_cb, 250 * LWS_US_PER_MS); - } lwsac_free(&a.ac); break; } @@ -679,19 +621,8 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) goto soft_error; r = sais_plat_reset(m->vhd, pr->event_uuid, pr->platform); - if (r == SAI_DB_RESULT_BUSY) { - sai_sul_retry_ctx_t *ctx = malloc(sizeof(*ctx)); - if (!ctx) - break; - - ctx->vhd = m->vhd; - lws_strncpy(ctx->uuid, pr->event_uuid, sizeof(ctx->uuid)); - lws_strncpy(ctx->platform, pr->platform, sizeof(ctx->platform)); - ctx->retries = 10; - ctx->op = SAIS_WS_WEBSRV_RX_PLATRESET; - lws_sul_schedule(m->vhd->context, 0, &ctx->sul, - sais_websrv_retry_cb, 250 * LWS_US_PER_MS); - } + if (r) + lwsl_ss_err(m->ss, "platreset failed"); lwsac_free(&a.ac); break; } @@ -709,18 +640,8 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) lwsl_notice("%s: eventdelete %s\n", __func__, ei->event_hash); r = sais_event_delete(m->vhd, ei->event_hash); - if (r == SAI_DB_RESULT_BUSY) { - sai_sul_retry_ctx_t *ctx = malloc(sizeof(*ctx)); - if (!ctx) - break; - - ctx->vhd = m->vhd; - lws_strncpy(ctx->uuid, ei->event_hash, sizeof(ctx->uuid)); - ctx->retries = 10; - ctx->op = SAIS_WS_WEBSRV_RX_EVENTDELETE; - lws_sul_schedule(m->vhd->context, 0, &ctx->sul, - sais_websrv_retry_cb, 250 * LWS_US_PER_MS); - } + if (r) + lwsl_ss_err(m->ss, "event delete failed"); lwsac_free(&a.ac); break; } diff --git a/src/web/w-comms.c b/src/web/w-comms.c index 8ce280f..4427885 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -441,7 +441,7 @@ saiw_event_db_close_all_now(struct vhd *vhd) } static int -callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, +w_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, void *in, size_t len) { struct vhd *vhd = (struct vhd *)lws_protocol_vh_priv_get( @@ -595,7 +595,12 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, "@com.warmcat.sai-websrv", 23)) lwsl_warn("%s: unable to set metadata\n", __func__); - return lws_ss_client_connect(vhd->h_ss_websrv) ? -1 : 0; + r = lws_ss_client_connect(vhd->h_ss_websrv) ? -1 : 0; + + if (r) + lwsl_wsi_err(wsi, "client connect for web -> srv failed"); + + return r; case LWS_CALLBACK_PROTOCOL_DESTROY: saiw_event_db_close_all_now(vhd); @@ -1026,8 +1031,10 @@ clean_spa: * Returning non-zero rejects it. */ if (n >= 8 && !strncmp((const char *)buf + n - 8, - "/builder", 8)) + "/builder", 8)) { + lwsl_wsi_err(wsi, "Terminating unexpected sai-web conn to /builder"); return 1; /* Reject builder connections */ + } return 0; @@ -1110,14 +1117,19 @@ clean_spa: if (lws_hdr_total_length(wsi, WSI_TOKEN_GET_URI)) { if (lws_hdr_copy(wsi, (char *)start, 64, - WSI_TOKEN_GET_URI) < 0) + WSI_TOKEN_GET_URI) < 0) { + lwsl_wsi_err(wsi, "URI too long"); return -1; + } } #if defined(LWS_ROLE_H2) else if (lws_hdr_copy(wsi, (char *)start, 64, - WSI_TOKEN_HTTP_COLON_PATH) < 0) + WSI_TOKEN_HTTP_COLON_PATH) < 0) { + lwsl_wsi_err(wsi, "path too long"); + return -1; + } #endif if (!memcmp((char *)start, "/sai", 4)) @@ -1201,7 +1213,6 @@ clean_spa: lws_dll2_foreach_safe(&pss->sched, NULL, saiw_sched_destroy); lwsac_free(&pss->logs_ac); - break; case LWS_CALLBACK_RECEIVE: @@ -1213,8 +1224,11 @@ clean_spa: /* * Browser UI sent us something on websockets */ - if (saiw_ws_json_rx_browser(vhd, pss, in, len)) + if (saiw_ws_json_rx_browser(vhd, pss, in, len)) { + lwsl_wsi_err(wsi, "Closing because saiw_ws_json_rx_browser returned it"); + return -1; + } break; @@ -1238,14 +1252,19 @@ passthru: return lws_callback_http_dummy(wsi, reason, user, in, len); bail: + lwsl_wsi_err(wsi, "Closing on bail"); + return 1; try_to_reuse: - if (lws_http_transaction_completed(wsi)) + if (lws_http_transaction_completed(wsi)) { + lwsl_wsi_err(wsi, "Closing because transaction_completed said so"); + return -1; + } return 0; } const struct lws_protocols protocol_ws = - { "com-warmcat-sai", callback_ws, sizeof(struct pss), 0 }; + { "com-warmcat-sai", w_callback_ws, sizeof(struct pss), 0 }; diff --git a/src/web/w-private.h b/src/web/w-private.h index 344396a..3380522 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -75,7 +75,9 @@ typedef struct { } sai_notification_t; typedef enum { - WSS_IDLE, + WSS_IDLE1, + WSS_IDLE2, + WSS_IDLE3, WSS_PREPARE_OVERVIEW, WSS_SEND_OVERVIEW, WSS_PREPARE_BUILDER_SUMMARY, diff --git a/src/web/w-websrv.c b/src/web/w-websrv.c index 5601a4e..1e16c29 100644 --- a/src/web/w-websrv.c +++ b/src/web/w-websrv.c @@ -277,6 +277,9 @@ saiw_websrv_queue_tx(struct lws_ss_handle *h, void *buf, size_t len) { saiw_websrv_t *m = (saiw_websrv_t *)lws_ss_to_user_object(h); + lwsl_ss_notice(h, "sai-web: Queuing from browser -> sai-server"); + lwsl_hexdump_notice(buf, len); + if (lws_buflist_append_segment(&m->bltx, buf, len) < 0) return 1; diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 743d3ce..68775d1 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -681,7 +681,7 @@ again: sch = NULL; switch (pss->send_state) { - case WSS_IDLE: + case WSS_IDLE1: /* * Anything from a task log he's subscribed to? @@ -701,7 +701,6 @@ again: */ if (pss->log_cache_index == pss->log_cache_size) { - // lws_usec_t tim; int sr; sai_task_uuid_to_event_uuid(event_uuid, @@ -726,8 +725,6 @@ again: return 0; } - // tim = lws_now_usecs(); - sr = lws_struct_sq3_deserialize(pdb, esc, "uid,timestamp ", lsm_schema_sq3_map_log, @@ -745,9 +742,6 @@ again: pss->log_cache_index = 0; pss->log_cache_size = (int)pss->logs_owner.count; - - // lwsl_wsi_notice(pss->wsi, "fetched %d logs in %dus", pss->log_cache_size, (int)(lws_now_usecs() - tim)); - } if (pss->log_cache_index < pss->log_cache_size) { @@ -793,45 +787,61 @@ again: } /* - * ...then do we have anything on the scheduled ll for this pss? + * Stay in this state if we're in the middle of a + * multi-fragment message */ + if (lws_ws_sending_multifragment(pss->wsi)) + return 0; - if (!sch) { - /* - * Send anything waiting on broadcast_raw buflist first - */ + /* fallthru */ - if (pss->raw_tx) { - char som, eom, rb[4096]; - int used, *pi = (int *)rb; + case WSS_IDLE2: - used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, - sizeof(rb), &som, &eom); - if (!used) - return 0; + pss->send_state = WSS_IDLE2; - // lwsl_wsi_notice(pss->wsi, "writing %d bytes flags 0x%x: '%.*s'", - // (int)(used - (int)sizeof(int)), (int)*pi, - // (int)(used - (int)sizeof(int)), rb + sizeof(int)); + /* + * Send anything waiting on broadcast_raw buflist first + */ - if (lws_write(pss->wsi, (uint8_t *)rb + sizeof(int), - (size_t)used - sizeof(int), - (enum lws_write_protocol)*pi) < 0) { - lwsl_wsi_err(pss->wsi, "attempt to write %d failed", (int)used - (int)sizeof(int)); + if (pss->raw_tx) { + char som, eom, rb[4096]; + int used, *pi = (int *)rb; - return -1; - } + used = lws_buflist_fragment_use(&pss->raw_tx, (uint8_t *)rb, + sizeof(rb), &som, &eom); + if (!used) + return 0; - lws_callback_on_writable(pss->wsi); + // lwsl_wsi_notice(pss->wsi, "writing %d bytes flags 0x%x: '%.*s'", + // (int)(used - (int)sizeof(int)), (int)*pi, + // (int)(used - (int)sizeof(int)), rb + sizeof(int)); - return 0; + if (lws_write(pss->wsi, (uint8_t *)rb + sizeof(int), + (size_t)used - sizeof(int), + (enum lws_write_protocol)*pi) < 0) { + lwsl_wsi_err(pss->wsi, "attempt to write %d failed", (int)used - (int)sizeof(int)); + + return -1; } + lws_callback_on_writable(pss->wsi); + + if (!lws_ws_sending_multifragment(pss->wsi)) + pss->send_state = WSS_IDLE1; - /* ... nope... */ return 0; } + /* + * Stay in this state if we're in the middle of a + * multi-fragment message, otherwise do whatever the + * sch proposes + */ + + if (lws_ws_sending_multifragment(pss->wsi) || + !sch) + return 0; + pss->toggle_favour_sch = 0; pss->send_state = sch->action; goto again; @@ -853,7 +863,7 @@ again: &sch->ac, 0, -8)) { lwsl_notice("%s: OVERVIEW 2 failed\n", __func__); - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; saiw_dealloc_sched(sch); return 0; @@ -955,7 +965,7 @@ again: lws_struct_json_serialize_destroy(&js); switch (n) { case LSJS_RESULT_ERROR: - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; saiw_dealloc_sched(sch); lwsl_err("%s: json ser error\n", __func__); return 1; @@ -970,7 +980,7 @@ again: } } if (!any) { - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; saiw_dealloc_sched(sch); return 0; @@ -1066,7 +1076,7 @@ enum_tasks: so_finish: p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), "]}"); - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; endo = 1; break; @@ -1145,7 +1155,7 @@ so_finish: switch (lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w)) { case LSJS_RESULT_ERROR: lws_struct_json_serialize_destroy(&js); - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; saiw_dealloc_sched(sch); return 1; case LSJS_RESULT_FINISH: @@ -1257,7 +1267,7 @@ b_finish: p += w; if (!lws_ptr_diff(p, start)) { saiw_dealloc_sched(sch); - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; lwsl_notice("%s: taskinfo: empty json\n", __func__); return 0; } @@ -1314,7 +1324,7 @@ send_it: if (lg || endo || - (pss->send_state == WSS_IDLE && sch) || + (pss->send_state == WSS_IDLE1 && sch) || (pss->send_state != WSS_SEND_ARTIFACT_INFO && sch && !sch->walk) || (pss->send_state == WSS_SEND_ARTIFACT_INFO && (!sch || !sch->owner.head))) { @@ -1329,7 +1339,7 @@ send_it: pss->sub_task_uuid); } - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; saiw_dealloc_sched(sch); } @@ -1342,7 +1352,7 @@ send_it: return 0; no_sch: - pss->send_state = WSS_IDLE; + pss->send_state = WSS_IDLE1; return 0; }
Page fetched 0s ago, creation time: 13ms (vhost etag hits: 0%, cache hits: 0%)